diff --git a/src/ClusterBootstrap/services/jobmanager/jobmanager.yaml b/src/ClusterBootstrap/services/jobmanager/jobmanager.yaml index caa7ae9d6..54ccf8c67 100755 --- a/src/ClusterBootstrap/services/jobmanager/jobmanager.yaml +++ b/src/ClusterBootstrap/services/jobmanager/jobmanager.yaml @@ -13,6 +13,9 @@ spec: labels: jobmanager-node: pod app: jobmanager + annotations: + prometheus.io/scrape: "true" + prometheus.io/path: "/metrics" spec: {% if cnf["dnsPolicy"] %} dnsPolicy: {{cnf["dnsPolicy"]}} @@ -39,7 +42,40 @@ spec: - mountPath: {{cnf["storage-mount-path"]}}/jobfiles name: dlwsdatajobfiles - mountPath: /var/log/dlworkspace - name: log + name: log + ports: + - containerPort: 9200 + hostPort: 9200 + name: job-mgr + protocol: TCP + - containerPort: 9201 + hostPort: 9201 + name: user-mgr + protocol: TCP + - containerPort: 9202 + hostPort: 9202 + name: node-mgr + protocol: TCP + - containerPort: 9203 + hostPort: 9203 + name: joblog-mgr + protocol: TCP + - containerPort: 9204 + hostPort: 9204 + name: cmd-mgr + protocol: TCP + - containerPort: 9205 + hostPort: 9205 + name: endpoint-mgr + protocol: TCP + readinessProbe: + failureThreshold: 3 + initialDelaySeconds: 3 + periodSeconds: 30 + successThreshold: 1 + tcpSocket: + port: 9200 + timeoutSeconds: 10 volumes: - name: certs hostPath: diff --git a/src/ClusterBootstrap/services/restfulapi/restfulapi.yaml b/src/ClusterBootstrap/services/restfulapi/restfulapi.yaml index e3d089c83..d1e5a6ce4 100755 --- a/src/ClusterBootstrap/services/restfulapi/restfulapi.yaml +++ b/src/ClusterBootstrap/services/restfulapi/restfulapi.yaml @@ -4,7 +4,7 @@ metadata: name: restfulapi namespace: default labels: - run: dlwsrestfulapi + run: dlwsrestfulapi spec: selector: matchLabels: @@ -15,13 +15,17 @@ spec: labels: restfulapi-node: pod app: restfulapi + annotations: + prometheus.io/scrape: "true" + prometheus.io/path: "/metrics" + prometheus.io/port: "5000" spec: - {% if cnf["dnsPolicy"] %} + {% if cnf["dnsPolicy"] %} dnsPolicy: {{cnf["dnsPolicy"]}} {% endif %} nodeSelector: restfulapi: active - hostNetwork: true + hostNetwork: true containers: - name: restfulapi image: {{cnf["worker-dockerregistry"]}}{{cnf["dockerprefix"]}}{{cnf["restfulapi"]}}:{{cnf["dockertag"]}} @@ -31,6 +35,10 @@ spec: name: apiconfig - mountPath: /var/log/apache2 name: log + ports: + - containerPort: 5000 + hostPort: 5000 + name: main {% if False %} {% for volume in cnf["mountpoints"] %} {% if cnf["mountpoints"][volume]["mountpoints"] is string and cnf["mountpoints"][volume]["mountpoints"]!="" %} @@ -42,7 +50,7 @@ spec: name: {{mp}} {% endfor %} {% endif %} - {% endfor %} + {% endfor %} {% endif %} volumes: - name: apiconfig @@ -60,14 +68,14 @@ spec: {% else %} {% for mp in cnf["mountpoints"][volume]["mountpoints"] %} - name: {{mp}} - hostPath: + hostPath: path: {{cnf["storage-mount-path"]}}/{{mp}} {% endfor %} {% endif %} - {% endfor %} + {% endfor %} {% endif %} tolerations: - key: CriticalAddonsOnly operator: Exists - key: node-role.kubernetes.io/master - effect: NoSchedule + effect: NoSchedule diff --git a/src/ClusterManager/cluster_manager.py b/src/ClusterManager/cluster_manager.py index 845bd1662..2fdab5851 100755 --- a/src/ClusterManager/cluster_manager.py +++ b/src/ClusterManager/cluster_manager.py @@ -5,10 +5,37 @@ import logging.config import sys import time +import argparse +import threading + +from prometheus_client.twisted import MetricsResource + +from twisted.web.server import Site +from twisted.web.resource import Resource +from twisted.internet import reactor logger = logging.getLogger(__name__) +class HealthResource(Resource): + def render_GET(self, request): + request.setHeader("Content-Type", "text/html; charset=utf-8") + return "Ok".encode("utf-8") + +def exporter_thread(port): + root = Resource() + root.putChild(b"metrics", MetricsResource()) + root.putChild(b"healthz", HealthResource()) + factory = Site(root) + reactor.listenTCP(port, factory) + reactor.run(installSignalHandlers=False) + +def setup_exporter_thread(port): + t = threading.Thread(target=exporter_thread, args=(port,), + name="exporter") + t.start() + return t + def create_log(logdir="/var/log/dlworkspace"): if not os.path.exists(logdir): os.system("mkdir -p " + logdir) @@ -18,17 +45,17 @@ def create_log(logdir="/var/log/dlworkspace"): logging.config.dictConfig(logging_config) -def Run(): +def Run(args): create_log() cwd = os.path.dirname(__file__) cmds = [ - ["python", os.path.join(cwd, "job_manager.py")], - ["python", os.path.join(cwd, "user_manager.py")], - ["python", os.path.join(cwd, "node_manager.py")], - ["python", os.path.join(cwd, "joblog_manager.py")], - ["python", os.path.join(cwd, "command_manager.py")], - ["python", os.path.join(cwd, "endpoint_manager.py")], + ["python", os.path.join(cwd, "job_manager.py"), "--port", str(args.j)], + ["python", os.path.join(cwd, "user_manager.py"), "--port", str(args.u)], + ["python", os.path.join(cwd, "node_manager.py"), "--port", str(args.n)], + ["python", os.path.join(cwd, "joblog_manager.py"), "--port", str(args.l)], + ["python", os.path.join(cwd, "command_manager.py"), "--port", str(args.c)], + ["python", os.path.join(cwd, "endpoint_manager.py"), "--port", str(args.e)], ] FNULL = open(os.devnull, "w") @@ -49,4 +76,13 @@ def Run(): if __name__ == "__main__": - sys.exit(Run()) + parser = argparse.ArgumentParser() + parser.add_argument("-j", help="port of job_manager", type=int, default=9200) + parser.add_argument("-u", help="port of user_manager", type=int, default=9201) + parser.add_argument("-n", help="port of node_manager", type=int, default=9202) + parser.add_argument("-l", help="port of joblog_manager", type=int, default=9203) + parser.add_argument("-c", help="port of command_manager", type=int, default=9204) + parser.add_argument("-e", help="port of endpoint_manager", type=int, default=9205) + args = parser.parse_args() + + sys.exit(Run(args)) diff --git a/src/ClusterManager/command_manager.py b/src/ClusterManager/command_manager.py index 9ad9b5956..fd117835b 100755 --- a/src/ClusterManager/command_manager.py +++ b/src/ClusterManager/command_manager.py @@ -31,6 +31,8 @@ import logging +from cluster_manager import setup_exporter_thread + logger = logging.getLogger(__name__) def RunCommand(command): @@ -66,4 +68,9 @@ def Run(): time.sleep(1) if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9204) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/ClusterManager/endpoint_manager.py b/src/ClusterManager/endpoint_manager.py index 68c961395..d95931969 100755 --- a/src/ClusterManager/endpoint_manager.py +++ b/src/ClusterManager/endpoint_manager.py @@ -1,6 +1,4 @@ -from config import config, GetStoragePath, GetWorkPath -import k8sUtils -from DataHandler import DataHandler + import json import os import time @@ -12,8 +10,16 @@ import random import re import logging +import yaml +import logging.config + +import argparse +from cluster_manager import setup_exporter_thread sys.path.append(os.path.join(os.path.dirname(os.path.abspath(__file__)), "../utils")) +import k8sUtils +from config import config, GetStoragePath, GetWorkPath +from DataHandler import DataHandler logger = logging.getLogger(__name__) @@ -246,4 +252,9 @@ def Run(): time.sleep(1) if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9205) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/ClusterManager/job_deployer.py b/src/ClusterManager/job_deployer.py index 5d2e22e12..e7cf76839 100644 --- a/src/ClusterManager/job_deployer.py +++ b/src/ClusterManager/job_deployer.py @@ -1,6 +1,7 @@ import yaml import os import logging +import logging.config from kubernetes import client, config from kubernetes.client.rest import ApiException from kubernetes.stream import stream @@ -98,6 +99,7 @@ def get_pods(self, field_selector="", label_selector=""): field_selector=field_selector, label_selector=label_selector, ) + logging.debug("Get pods: {}".format(api_response)) return api_response.items def get_services_by_label(self, label_selector): @@ -125,29 +127,33 @@ def delete_job(self, job_id): def pod_exec(self, pod_name, exec_command, timeout=60): """work as the command (with timeout): kubectl exec 'pod_name' 'exec_command'""" - logging.info("Exec on pod {}: {}".format(pod_name, exec_command)) - client = stream( - self.v1.connect_get_namespaced_pod_exec, - name=pod_name, - namespace=self.namespace, - command=exec_command, - stderr=True, - stdin=False, - stdout=True, - tty=False, - _preload_content=False, - ) - client.run_forever(timeout=timeout) - - err = yaml.full_load(client.read_channel(ERROR_CHANNEL)) - if err is None: - return [-1, "Timeout"] - - if err["status"] == "Success": - status_code = 0 - else: - logging.warning("Exec on pod {} failed. cmd: {}, err: {}.".format(pod_name, exec_command, err)) - status_code = int(err["details"]["causes"][0]["message"]) - output = client.read_all() - logging.info("Exec on pod {}, status: {}, cmd: {}, output: {}".format(pod_name, status_code, exec_command, output)) - return [status_code, output] + try: + logging.info("Exec on pod {}: {}".format(pod_name, exec_command)) + client = stream( + self.v1.connect_get_namespaced_pod_exec, + name=pod_name, + namespace=self.namespace, + command=exec_command, + stderr=True, + stdin=False, + stdout=True, + tty=False, + _preload_content=False, + ) + client.run_forever(timeout=timeout) + + err = yaml.full_load(client.read_channel(ERROR_CHANNEL)) + if err is None: + return [-1, "Timeout"] + + if err["status"] == "Success": + status_code = 0 + else: + logging.debug("Exec on pod {} failed. cmd: {}, err: {}.".format(pod_name, exec_command, err)) + status_code = int(err["details"]["causes"][0]["message"]) + output = client.read_all() + logging.info("Exec on pod {}, status: {}, cmd: {}, output: {}".format(pod_name, status_code, exec_command, output)) + return [status_code, output] + except ApiException as err: + logging.error("Exec on pod {} error. cmd: {}, err: {}.".format(pod_name, exec_command, err), exc_info=True) + return [-1, err.message] diff --git a/src/ClusterManager/job_manager.py b/src/ClusterManager/job_manager.py index 221852ae4..fdee7f47c 100755 --- a/src/ClusterManager/job_manager.py +++ b/src/ClusterManager/job_manager.py @@ -9,7 +9,6 @@ import copy import traceback - sys.path.append(os.path.join(os.path.dirname(os.path.abspath(__file__)),"../storage")) sys.path.append(os.path.join(os.path.dirname(os.path.abspath(__file__)),"../utils")) @@ -41,8 +40,24 @@ from job_deployer import JobDeployer from job_role import JobRole +from cluster_manager import setup_exporter_thread + + +def all_pods_not_existing(job_id): + job_deployer = JobDeployer() + job_roles = JobRole.get_job_roles(job_id) + statuses = [job_role.status() for job_role in job_roles] + logging.info("Job: {}, status: {}".format(job_id, statuses)) + return all([status == "NotFound" for status in statuses]) + def SubmitJob(job): + # check if existing any pod with label: run=job_id + assert("jobId" in job) + if not all_pods_not_existing(job["jobId"]): + logging.warning("Waiting until previously pods are cleaned up! Job {}".format(job["jobId"])) + return + ret = {} dataHandler = DataHandler() @@ -118,31 +133,31 @@ def SubmitJob(job): return ret -def KillJob(job, desiredState="killed"): +def KillJob(job_id, desiredState="killed"): dataHandler = DataHandler() - result, detail = k8sUtils.GetJobStatus(job["jobId"]) - dataHandler.UpdateJobTextField(job["jobId"], "jobStatusDetail", base64.b64encode(json.dumps(detail))) - logging.info("Killing job %s, with status %s, %s" % (job["jobId"], result, detail)) + result, detail = k8sUtils.GetJobStatus(job_id) + dataHandler.UpdateJobTextField(job_id, "jobStatusDetail", base64.b64encode(json.dumps(detail))) + logging.info("Killing job %s, with status %s, %s" % (job_id, result, detail)) job_deployer = JobDeployer() - errors = job_deployer.delete_job(job["jobId"]) + errors = job_deployer.delete_job(job_id) if len(errors) == 0: - dataHandler.UpdateJobTextField(job["jobId"], "jobStatus", desiredState) - dataHandler.UpdateJobTextField(job["jobId"], "lastUpdated", datetime.datetime.now().isoformat()) + dataHandler.UpdateJobTextField(job_id, "jobStatus", desiredState) + dataHandler.UpdateJobTextField(job_id, "lastUpdated", datetime.datetime.now().isoformat()) dataHandler.Close() return True else: - dataHandler.UpdateJobTextField(job["jobId"], "jobStatus", "error") - dataHandler.UpdateJobTextField(job["jobId"], "lastUpdated", datetime.datetime.now().isoformat()) + dataHandler.UpdateJobTextField(job_id, "jobStatus", "error") + dataHandler.UpdateJobTextField(job_id, "lastUpdated", datetime.datetime.now().isoformat()) dataHandler.Close() logging.error("Kill job failed with errors: {}".format(errors)) return False -def ApproveJob(job): +def ApproveJob(job_id): dataHandler = DataHandler() - dataHandler.UpdateJobTextField(job["jobId"], "jobStatus", "queued") + dataHandler.UpdateJobTextField(job_id, "jobStatus", "queued") dataHandler.Close() return True @@ -188,9 +203,14 @@ def UpdateJobStatus(job): if jobDescriptionPath is not None and os.path.isfile(jobDescriptionPath): k8sUtils.kubectl_delete(jobDescriptionPath) - elif result == "Unknown": + elif result == "Unknown" or result == "NotFound": if job["jobId"] not in UnusualJobs: + logging.warning("!!! Job status ---{}---, job: {}".format(result, job["jobId"])) UnusualJobs[job["jobId"]] = datetime.datetime.now() + # TODO + # 1) May need to reduce the timeout. + # It takes minutes before pod turns into "Unknown", we may don't need to wait so long. + # 2) If node resume before we resubmit the job, the job will end in status 'NotFound'. elif (datetime.datetime.now() - UnusualJobs[job["jobId"]]).seconds > 300: del UnusualJobs[job["jobId"]] retries = dataHandler.AddandGetJobRetries(job["jobId"]) @@ -202,13 +222,10 @@ def UpdateJobStatus(job): k8sUtils.kubectl_delete(jobDescriptionPath) else: logging.warning("Job %s fails in Kubernetes, delete and re-submit the job. Retries %d", job["jobId"], retries) - SubmitJob(job) - elif result.strip() == "PendingHostPort": - logging.warning("Cannot find host ports for job :%s, re-launch the job with different host ports ", job["jobId"]) - - SubmitJob(job) + KillJob(job["jobId"], "queued") + # SubmitJob(job) - if result.strip() != "Unknown" and job["jobId"] in UnusualJobs: + if result != "Unknown" and result != "NotFound" and job["jobId"] in UnusualJobs: del UnusualJobs[job["jobId"]] dataHandler.Close() @@ -219,6 +236,9 @@ def check_job_status(job_id): job_deployer = JobDeployer() job_roles = JobRole.get_job_roles(job_id) + if len(job_roles) < 1: + return "NotFound" + # role status in ["NotFound", "Pending", "Running", "Succeeded", "Failed", "Unknown"] # TODO ??? when ps/master role "Succeeded", return Succeeded for job_role in job_roles: @@ -327,12 +347,15 @@ def TakeJobActions(jobs): logging.info("TakeJobActions : global resources : %s" % (globalResInfo.CategoryToCountMap)) for sji in jobsInfo: - if sji["job"]["jobStatus"] == "queued" and sji["allowed"] == True: - SubmitJob(sji["job"]) - logging.info("TakeJobActions : submitting job : %s : %s : %s" % (sji["jobParams"]["jobName"], sji["jobParams"]["jobId"], sji["sortKey"])) - elif sji["jobParams"]["preemptionAllowed"] and (sji["job"]["jobStatus"] == "scheduling" or sji["job"]["jobStatus"] == "running") and sji["allowed"] == False: - KillJob(sji["job"], "queued") - logging.info("TakeJobActions : pre-empting job : %s : %s : %s" % (sji["jobParams"]["jobName"], sji["jobParams"]["jobId"], sji["sortKey"])) + try: + if sji["job"]["jobStatus"] == "queued" and sji["allowed"] == True: + SubmitJob(sji["job"]) + logging.info("TakeJobActions : submitting job : %s : %s : %s" % (sji["jobParams"]["jobName"], sji["jobParams"]["jobId"], sji["sortKey"])) + elif sji["jobParams"]["preemptionAllowed"] and (sji["job"]["jobStatus"] == "scheduling" or sji["job"]["jobStatus"] == "running") and sji["allowed"] == False: + KillJob(sji["job"]["jobId"], "queued") + logging.info("TakeJobActions : pre-empting job : %s : %s : %s" % (sji["jobParams"]["jobName"], sji["jobParams"]["jobId"], sji["sortKey"])) + except Exception as e: + logging.error("Process job failed {}".format(sji["job"]), exc_info=True) logging.info("TakeJobActions : job desired actions taken") @@ -341,7 +364,6 @@ def Run(): create_log() while True: - try: config["racks"] = k8sUtils.get_node_labels("rack") config["skus"] = k8sUtils.get_node_labels("sku") @@ -360,15 +382,15 @@ def Run(): try: logging.info("Processing job: %s, status: %s" % (job["jobId"], job["jobStatus"])) if job["jobStatus"] == "killing": - KillJob(job, "killed") + KillJob(job["jobId"], "killed") elif job["jobStatus"] == "pausing": - KillJob(job, "paused") + KillJob(job["jobId"], "paused") elif job["jobStatus"] == "scheduling" or job["jobStatus"] == "running": UpdateJobStatus(job) elif job["jobStatus"] == "unapproved": - ApproveJob(job) + ApproveJob(job["jobId"]) except Exception as e: - logging.info(e) + logging.info(e, exc_info=True) except Exception as e: logging.exception("process pending job failed") finally: @@ -380,4 +402,9 @@ def Run(): if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9200) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/ClusterManager/job_role.py b/src/ClusterManager/job_role.py index 7ceb00ecf..f4e61e9f8 100644 --- a/src/ClusterManager/job_role.py +++ b/src/ClusterManager/job_role.py @@ -1,3 +1,5 @@ +import logging +import logging.config from job_deployer import JobDeployer @@ -31,8 +33,10 @@ def status(self): CONTAINER_READY -> WORKER_READY -> JOB_READY (then the job finally in "Running" status.) """ # pod-phase: https://kubernetes.io/docs/concepts/workloads/pods/pod-lifecycle/#pod-phase + # node condition: https://kubernetes.io/docs/concepts/architecture/nodes/#condition deployer = JobDeployer() pods = deployer.get_pods(field_selector="metadata.name={}".format(self.pod_name)) + logging.debug("Pods: {}".format(pods)) if(len(pods) < 1): return "NotFound" @@ -40,8 +44,13 @@ def status(self): pod = pods[0] phase = pod.status.phase - # !!! Pod is runing, doesn't mean "Role" is ready and running. + # !!! Pod is running, doesn't mean "Role" is ready and running. if(phase == "Running"): + # Found that phase won't turn into "Unkonwn" even when we get 'unknown' from kubectl + if pod.status.reason == "NodeLost": + return "Unknown" + + # Check if the user command had been ran. if not self.isRoleReady(): return "Pending" diff --git a/src/ClusterManager/job_status.pdf b/src/ClusterManager/job_status.pdf new file mode 100644 index 000000000..c9756f120 Binary files /dev/null and b/src/ClusterManager/job_status.pdf differ diff --git a/src/ClusterManager/joblog_manager.py b/src/ClusterManager/joblog_manager.py index 5e55e74d1..bdc383e62 100755 --- a/src/ClusterManager/joblog_manager.py +++ b/src/ClusterManager/joblog_manager.py @@ -32,6 +32,8 @@ from config import config, GetStoragePath from DataHandler import DataHandler +from cluster_manager import setup_exporter_thread + logger = logging.getLogger(__name__) def create_log(logdir = '/var/log/dlworkspace'): @@ -159,4 +161,9 @@ def Run(): time.sleep(1) if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9203) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/ClusterManager/logging.yaml b/src/ClusterManager/logging.yaml index b276c6d8d..a486bc5aa 100755 --- a/src/ClusterManager/logging.yaml +++ b/src/ClusterManager/logging.yaml @@ -1,26 +1,27 @@ -version: 1 -formatters: - simple: - format: '%(asctime)s - %(name)s - %(levelname)s - %(message)s' -handlers: - console: - class: logging.StreamHandler - level: DEBUG - formatter: simple - stream: ext://sys.stdout - file: - class : logging.handlers.RotatingFileHandler - formatter: simple - filename: /var/log/dlworkspace/clustermanager.log - # roll over at 10MB - maxBytes: 10240000 - # At most 10 logging files - backupCount: 10 -loggers: - basic: - level: DEBUG - handlers: ['console','file'] - propagate: no -root: - level: DEBUG - handlers: ['console','file'] \ No newline at end of file +version: 1 +disable_existing_loggers: False +formatters: + simple: + format: '%(asctime)s - %(levelname)s - %(filename)s:%(lineno)d - %(message)s' +handlers: + console: + class: logging.StreamHandler + level: INFO + formatter: simple + stream: ext://sys.stdout + file: + class : logging.handlers.RotatingFileHandler + formatter: simple + filename: /var/log/dlworkspace/clustermanager.log + # roll over at 10MB + maxBytes: 10240000 + # At most 10 logging files + backupCount: 10 +loggers: + basic: + level: INFO + handlers: ['console','file'] + propagate: no +root: + level: INFO + handlers: ['console','file'] diff --git a/src/ClusterManager/node_manager.py b/src/ClusterManager/node_manager.py index 36c3b5b5c..0849a6038 100755 --- a/src/ClusterManager/node_manager.py +++ b/src/ClusterManager/node_manager.py @@ -39,6 +39,7 @@ from config import config from DataHandler import DataHandler +from cluster_manager import setup_exporter_thread def create_log(logdir = '/var/log/dlworkspace'): @@ -252,4 +253,9 @@ def Run(): time.sleep(30) if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9202) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/ClusterManager/requirements.txt b/src/ClusterManager/requirements.txt index 111149103..f1363b2db 100644 --- a/src/ClusterManager/requirements.txt +++ b/src/ClusterManager/requirements.txt @@ -1,3 +1,5 @@ marshmallow==2.19.5 kubernetes==9.0.0 PyYAML>=5.1.1 +prometheus-client==0.7.1 +twisted==19.2.1 diff --git a/src/ClusterManager/user_manager.py b/src/ClusterManager/user_manager.py index 13968f365..fc89f0b18 100755 --- a/src/ClusterManager/user_manager.py +++ b/src/ClusterManager/user_manager.py @@ -34,6 +34,7 @@ from config import config from DataHandler import DataHandler +from cluster_manager import setup_exporter_thread def create_log(logdir = '/var/log/dlworkspace'): @@ -91,4 +92,9 @@ def Run(): if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", help="port of exporter", type=int, default=9201) + args = parser.parse_args() + setup_exporter_thread(args.port) + Run() diff --git a/src/RestAPI/dlwsrestapi.py b/src/RestAPI/dlwsrestapi.py index df34e6600..77477c549 100755 --- a/src/RestAPI/dlwsrestapi.py +++ b/src/RestAPI/dlwsrestapi.py @@ -2,7 +2,7 @@ import json import os -from flask import Flask +from flask import Flask, Response from flask_restful import reqparse, abort, Api, Resource from flask import request, jsonify import base64 @@ -28,6 +28,10 @@ import traceback import threading +import prometheus_client + +CONTENT_TYPE_LATEST = str("text/plain; version=0.0.4; charset=utf-8") + dir_path = os.path.dirname(os.path.realpath(__file__)) with open(os.path.join(dir_path, 'logging.yaml'), 'r') as f: logging_config = yaml.load(f) @@ -1205,6 +1209,10 @@ def endpoint_exist(endpoint_id): ## api.add_resource(Endpoint, '/endpoints') +@app.route("/metrics") +def metrics(): + return Response(prometheus_client.generate_latest(), mimetype=CONTENT_TYPE_LATEST) + if __name__ == '__main__': app.run(debug=False,host="0.0.0.0",threaded=True) diff --git a/src/RestAPI/logging.yaml b/src/RestAPI/logging.yaml index d42108867..fea884c5c 100755 --- a/src/RestAPI/logging.yaml +++ b/src/RestAPI/logging.yaml @@ -1,26 +1,27 @@ -version: 1 -formatters: - simple: - format: '%(asctime)s - %(name)s - %(levelname)s - %(message)s' -handlers: - console: - class: logging.StreamHandler - level: DEBUG - formatter: simple - stream: ext://sys.stdout - file: - class : logging.handlers.RotatingFileHandler - formatter: simple - filename: /var/log/apache2/restfulapi.log - # roll over at 10MB - maxBytes: 10240000 - # At most 10 logging files - backupCount: 10 -loggers: - basic: - level: DEBUG - handlers: ['console', 'file'] - propagate: no -root: - level: DEBUG - handlers: ['console', 'file'] +version: 1 +disable_existing_loggers: False +formatters: + simple: + format: '%(asctime)s - %(levelname)s - %(filename)s:%(lineno)d - %(message)s' +handlers: + console: + class: logging.StreamHandler + level: INFO + formatter: simple + stream: ext://sys.stdout + file: + class : logging.handlers.RotatingFileHandler + formatter: simple + filename: /var/log/apache2/restfulapi.log + # roll over at 10MB + maxBytes: 10240000 + # At most 10 logging files + backupCount: 10 +loggers: + basic: + level: INFO + handlers: ['console', 'file'] + propagate: no +root: + level: INFO + handlers: ['console', 'file'] diff --git a/src/utils/MySQLDataHandler.py b/src/utils/MySQLDataHandler.py index 380138c79..7ced12ca8 100755 --- a/src/utils/MySQLDataHandler.py +++ b/src/utils/MySQLDataHandler.py @@ -1,9 +1,9 @@ -# from config import config import mysql.connector import json import base64 import os import logging +import functools import timeit @@ -12,12 +12,35 @@ from config import config from config import global_vars +from prometheus_client import Histogram + logger = logging.getLogger(__name__) +data_handler_fn_histogram = Histogram("datahandler_fn_latency_seconds", + "latency for executing data handler function (seconds)", + buckets=(.05, .075, .1, .25, .5, .75, 1.0, 2.5, 5.0, + 7.5, 10.0, 12.5, 15.0, 17.5, 20.0, float("inf")), + labelnames=("fn_name",)) + +db_connect_histogram = Histogram("db_connect_latency_seconds", + "latency for connecting to db (seconds)", + buckets=(.05, .075, .1, .25, .5, .75, 1.0, 2.5, 5.0, 7.5, float("inf"))) + -class DataHandler: +def record(fn): + @functools.wraps(fn) + def wrapped(*args, **kwargs): + start = timeit.default_timer() + try: + return fn(*args, **kwargs) + finally: + elapsed = timeit.default_timer() - start + logger.info("DataHandler: %s, time elapsed %.2fs", fn.__name__, elapsed) + data_handler_fn_histogram.labels(fn.__name__).observe(elapsed) + return wrapped +class DataHandler(object): def __init__(self): start_time = timeit.default_timer() self.database = "DLWSCluster-%s" % config["clusterId"] @@ -32,9 +55,9 @@ def __init__(self): username = config["mysql"]["username"] password = config["mysql"]["password"] - - self.conn = mysql.connector.connect(user=username, password=password, - host=server, database=self.database) + with db_connect_histogram.time(): + self.conn = mysql.connector.connect(user=username, password=password, + host=server, database=self.database) self.CreateDatabase() self.CreateTable() @@ -43,8 +66,6 @@ def __init__(self): elapsed = timeit.default_timer() - start_time logger.info("DataHandler initialization, time elapsed %f s", elapsed) - - def CreateDatabase(self): if "initSQLDB" not in global_vars or not global_vars["initSQLDB"]: logger.info("===========init SQL database===============") @@ -218,41 +239,34 @@ def CreateTable(self): self.conn.commit() cursor.close() - + @record def AddStorage(self, vcName, url, storageType, metadata, defaultMountPath): try: - start_time = timeit.default_timer() sql = "INSERT INTO `"+self.storagetablename+"` (storageType, url, metadata, vcName, defaultMountPath) VALUES (%s,%s,%s,%s,%s)" cursor = self.conn.cursor() cursor.execute(sql, (storageType, url, metadata, vcName, defaultMountPath)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: AddStorage to DB: url : %s, vcName: %s , time elapsed %f s", url, vcName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False - + @record def DeleteStorage(self, vcName, url): try: - start_time = timeit.default_timer() sql = "DELETE FROM `%s` WHERE url = '%s' and vcName = '%s'" % (self.storagetablename, url, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: DeleteStorage: url:%s, vcName:%s, time elapsed %f s", url, vcName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False - + @record def ListStorages(self, vcName): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `storageType`,`url`,`metadata`,`vcName`,`defaultMountPath` FROM `%s` WHERE vcName = '%s' " % (self.storagetablename, vcName) ret = [] @@ -271,45 +285,39 @@ def ListStorages(self, vcName): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: ListStorages time elapsed %f s", elapsed) return ret + @record def UpdateStorage(self, vcName, url, storageType, metadata, defaultMountPath): try: - start_time = timeit.default_timer() sql = """update `%s` set storageType = '%s', metadata = '%s', defaultMountPath = '%s' where vcName = '%s' and url = '%s' """ % (self.storagetablename, storageType, metadata, defaultMountPath, vcName, url) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: UpdateStorage: vcName: %s, url: %s, time elapsed %f s", vcName, url, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def AddVC(self, vcName, quota, metadata): try: - start_time = timeit.default_timer() sql = "INSERT INTO `"+self.vctablename+"` (vcName, quota, metadata) VALUES (%s,%s,%s)" cursor = self.conn.cursor() cursor.execute(sql, (vcName, quota, metadata)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: AddVC to DB: vcName: %s , time elapsed %f s", vcName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def ListVCs(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `vcName`,`quota`,`metadata` FROM `%s`" % (self.vctablename) ret = [] @@ -326,45 +334,39 @@ def ListVCs(self): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: ListVCs time elapsed %f s", elapsed) return ret + @record def DeleteVC(self, vcName): try: - start_time = timeit.default_timer() sql = "DELETE FROM `%s` WHERE vcName = '%s'" % (self.vctablename, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: DeleteVC: vcName: %s , time elapsed %f s", vcName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def UpdateVC(self, vcName, quota, metadata): try: - start_time = timeit.default_timer() sql = """update `%s` set quota = '%s', metadata = '%s' where vcName = '%s' """ % (self.vctablename, quota, metadata, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: UpdateVC: vcName: %s , time elapsed %f s", vcName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetIdentityInfo(self, identityName): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `identityName`,`uid`,`gid`,`groups` FROM `%s` where `identityName` = '%s'" % (self.identitytablename, identityName) ret = [] @@ -379,17 +381,14 @@ def GetIdentityInfo(self, identityName): ret.append(record) except Exception as e: logger.error('GetIdentityInfo Exception: %s', str(e)) - pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: GetIdentityInfo time elapsed %f s", elapsed) return ret + @record def UpdateIdentityInfo(self, identityName, uid, gid, groups): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() if (isinstance(groups, list)): @@ -404,16 +403,14 @@ def UpdateIdentityInfo(self, identityName, uid, gid, groups): self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: UpdateIdentityInfo %s to database , time elapsed %f s", identityName, elapsed) return True except Exception as e: logger.error('UpdateIdentityInfo Exception: %s', str(e)) return False + @record def GetAceCount(self, identityName, resource): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM `%s` where `identityName` = '%s' and `resource` = '%s'" % (self.acltablename,identityName, resource) cursor.execute(query) @@ -422,14 +419,12 @@ def GetAceCount(self, identityName, resource): ret = c[0] self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: GetAceCount time elapsed %f s", elapsed) return ret + @record def UpdateAce(self, identityName, identityId, resource, permissions, isDeny): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() existingAceCount = self.GetAceCount(identityName, resource) logger.info(existingAceCount) @@ -443,34 +438,30 @@ def UpdateAce(self, identityName, identityId, resource, permissions, isDeny): self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: UpdateAce %s - %s to database , time elapsed %f s", identityName, resource, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def UpdateAclIdentityId(self, identityName, identityId): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() sql = """update `%s` set identityId = '%s' where `identityName` = '%s' """ % (self.acltablename, identityId, identityName) cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: UpdateAclIdentityId %s - %s to database , time elapsed %f s", identityName, identityId, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def DeleteResourceAcl(self, resource): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() sql = "DELETE FROM `%s` WHERE `resource` = '%s'" % (self.acltablename, resource) @@ -478,17 +469,15 @@ def DeleteResourceAcl(self, resource): self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: DeleteResourceAcl %s, time elapsed %f s", resource, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def DeleteAce(self, identityName, resource): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() sql = "DELETE FROM `%s` WHERE `identityName` = '%s' and `resource` = '%s'" % (self.acltablename, identityName, resource) @@ -496,16 +485,14 @@ def DeleteAce(self, identityName, resource): self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: DeleteAce %s : %s time elapsed %f s", resource, identityName, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetAcl(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `identityName`,`identityId`,`resource`,`permissions`,`isDeny` FROM `%s`" % (self.acltablename) ret = [] @@ -524,13 +511,11 @@ def GetAcl(self): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: GetAcl time elapsed %f s", elapsed) return ret + @record def GetResourceAcl(self, resource): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `identityName`,`identityId`,`resource`,`permissions`,`isDeny` FROM `%s` where `resource` = '%s'" % (self.acltablename, resource) ret = [] @@ -548,30 +533,26 @@ def GetResourceAcl(self, resource): logger.error('Exception: %s', str(e)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: GetResourceAcl time elapsed %f s", elapsed) return ret + @record def AddJob(self, jobParams): try: - start_time = timeit.default_timer() sql = "INSERT INTO `"+self.jobtablename+"` (jobId, familyToken, isParent, jobName, userName, vcName, jobType,jobParams ) VALUES (%s,%s,%s,%s,%s,%s,%s,%s)" cursor = self.conn.cursor() jobParam = base64.b64encode(json.dumps(jobParams)) cursor.execute(sql, (jobParams["jobId"], jobParams["familyToken"], jobParams["isParent"], jobParams["jobName"], jobParams["userName"], jobParams["vcName"], jobParams["jobType"],jobParam)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: added job %s to database, time elapsed %f s", jobParams["jobId"], elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or")): - start_time = timeit.default_timer() ret = [] cursor = self.conn.cursor() try: @@ -591,13 +572,12 @@ def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or" if num is not None: query += " limit %s " % str(num) - start_time1 = timeit.default_timer() cursor.execute(query) - elapsed1 = timeit.default_timer() - start_time1 - start_time2 = timeit.default_timer() + + fetch_start_time = timeit.default_timer() data = cursor.fetchall() - elapsed2 = timeit.default_timer() - start_time2 - logger.info("(fetchall time: %f)", elapsed2) + fetch_elapsed = timeit.default_timer() - fetch_start_time + logger.info("(fetchall time: %f)", fetch_elapsed) for (jobId,jobName,userName, vcName, jobStatus,jobStatusDetail, jobType, jobDescriptionPath, jobDescription, jobTime, endpoints, jobParams,errorMsg, jobMeta) in data: record = {} record["jobId"] = jobId @@ -619,12 +599,10 @@ def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or" logger.error('Exception: %s', str(e)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get job list for user %s , time elapsed %f s (SQL query time: %f)", userName, elapsed, elapsed1) return ret + @record def GetJob(self, **kwargs): - start_time = timeit.default_timer() valid_keys = ["jobId", "familyToken", "isParent", "jobName", "userName", "vcName", "jobStatus", "jobType", "jobTime"] if len(kwargs) != 1: return [] key, expected = kwargs.popitem() @@ -638,27 +616,23 @@ def GetJob(self, **kwargs): ret = [dict(zip(columns, row)) for row in cursor.fetchall()] self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get job details with query %s=%s , time elapsed %f s", key, expected, elapsed) return ret + @record def AddCommand(self, jobId, command): try: - start_time = timeit.default_timer() sql = "INSERT INTO `"+self.commandtablename+"` (jobId, command) VALUES (%s,%s)" cursor = self.conn.cursor() cursor.execute(sql, (jobId, command)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: add command to database, jobId: %s , time elapsed %f s", jobId, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetPendingCommands(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `id`, `jobId`, `command` FROM `%s` WHERE `status` = 'pending' order by `time`" % (self.commandtablename) cursor.execute(query) @@ -671,27 +645,23 @@ def GetPendingCommands(self): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get pending command , time elapsed %f s", elapsed) return ret + @record def FinishCommand(self, commandId): try: - start_time = timeit.default_timer() sql = """update `%s` set status = 'run' where `id` = '%s' """ % (self.commandtablename, commandId) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: set command %s as finished , time elapsed %f s", commandId, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetCommands(self, jobId): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `time`, `command`, `status`, `output` FROM `%s` WHERE `jobId` = '%s' order by `time`" % (self.commandtablename, jobId) cursor.execute(query) @@ -705,8 +675,6 @@ def GetCommands(self, jobId): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get command list for job %s , time elapsed %f s", jobId, elapsed) return ret def load_json(self, raw_str): @@ -717,9 +685,9 @@ def load_json(self, raw_str): except: return {} + @record def GetPendingEndpoints(self): try: - start_time = timeit.default_timer() jobs = self.GetJob(jobStatus="running") # [ {endpoint1:{},endpoint2:{}}, {endpoint3:{}, ... }, ... ] @@ -728,16 +696,14 @@ def GetPendingEndpoints(self): # endpoint["status"] == "pending" pendingEndpoints = {k: v for d in endpoints for k, v in d.items() if v["status"] == "pending"} - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get pending endpoints %d, time elapsed %f s", len(pendingEndpoints), elapsed) return pendingEndpoints except Exception as e: logger.exception("Query pending endpoints failed!") return {} + @record def GetDeadEndpoints(self): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() # TODO we need job["lastUpdated"] for filtering query = "SELECT `endpoints` FROM jobs WHERE `jobStatus` <> 'running' order by `jobTime` DESC" @@ -748,16 +714,14 @@ def GetDeadEndpoints(self): dead_endpoints.update(endpoint_list) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get dead endpoints %d , time elapsed %f s", len(dead_endpoints), elapsed) return dead_endpoints except Exception as e: logger.exception("Query dead endpoints failed!") return {} + @record def UpdateEndpoint(self, endpoint): try: - start_time = timeit.default_timer() job_id = endpoint["jobId"] job = self.GetJob(jobId=job_id)[0] job_endpoints = self.load_json(job["endpoints"]) @@ -770,15 +734,13 @@ def UpdateEndpoint(self, endpoint): cursor.execute(sql, (json.dumps(job_endpoints), job_id)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: update endpoints to database, endpointId: %s , time elapsed %f s", endpoint["id"], elapsed) return True except Exception as e: logger.exception("Update endpoints failed!") return False + @record def GetPendingJobs(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `jobId`,`jobName`,`userName`, `vcName`, `jobStatus`, `jobType`, `jobDescriptionPath`, `jobDescription`, `jobTime`, `endpoints`, `jobParams`,`errorMsg` ,`jobMeta` FROM `%s` where `jobStatus` <> 'error' and `jobStatus` <> 'failed' and `jobStatus` <> 'finished' and `jobStatus` <> 'killed' order by `jobTime` DESC" % (self.jobtablename) cursor.execute(query) @@ -801,44 +763,37 @@ def GetPendingJobs(self): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get pending jobs %d, time elapsed %f s", len(ret), elapsed) return ret - + @record def SetJobError(self, jobId, errorMsg): try: - start_time = timeit.default_timer() sql = """update `%s` set jobStatus = 'error', `errorMsg` = '%s' where `jobId` = '%s' """ % (self.jobtablename, errorMsg, jobId) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: set job %s error status in database, time elapsed %f s", jobId, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def UpdateJobTextField(self, jobId, field, value): try: - start_time = timeit.default_timer() sql = "update `%s` set `%s` = '%s' where `jobId` = '%s' " % (self.jobtablename, field, value, jobId) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: update job %s, field %s to %s, time elapsed %f s", jobId, field, value, elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetJobTextField(self, jobId, field): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `jobId`, `%s` FROM `%s` where `jobId` = '%s' " % (field, self.jobtablename, jobId) ret = None @@ -851,12 +806,10 @@ def GetJobTextField(self, jobId, field): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get filed %s of job %s , time elapsed %f s", field, jobId, elapsed) return ret + @record def AddandGetJobRetries(self, jobId): - start_time = timeit.default_timer() sql = """update `%s` set `retries` = `retries` + 1 where `jobId` = '%s' """ % (self.jobtablename, jobId) cursor = self.conn.cursor() cursor.execute(sql) @@ -872,29 +825,25 @@ def AddandGetJobRetries(self, jobId): ret = value self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get and update retries for job %s , time elapsed %f s", jobId, elapsed) return ret + @record def UpdateClusterStatus(self, clusterStatus): try: status = base64.b64encode(json.dumps(clusterStatus)) - start_time = timeit.default_timer() sql = "INSERT INTO `%s` (status) VALUES ('%s')" % (self.clusterstatustablename, status) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: update cluster status, time elapsed %f s", elapsed) return True except Exception as e: logger.error('Exception: %s', str(e)) return False + @record def GetClusterStatus(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `time`, `status` FROM `%s` order by `time` DESC limit 1" % (self.clusterstatustablename) ret = None @@ -908,13 +857,11 @@ def GetClusterStatus(self): logger.error('Exception: %s', str(e)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get cluster status , time elapsed %f s", elapsed) return ret, time + @record def GetUsers(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT `identityName`,`uid` FROM `%s`" % (self.identitytablename) ret = [] @@ -926,10 +873,9 @@ def GetUsers(self): logger.error('Exception: %s', str(e)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info("DataHandler: get users, time elapsed %f s", elapsed) return ret + @record def GetActiveJobsCount(self): cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM `%s` where `jobStatus` <> 'error' and `jobStatus` <> 'failed' and `jobStatus` <> 'finished' and `jobStatus` <> 'killed' " % (self.jobtablename) @@ -942,6 +888,7 @@ def GetActiveJobsCount(self): return ret + @record def GetALLJobsCount(self): cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM `%s`" % (self.jobtablename) diff --git a/src/utils/SQLDataHandler.py b/src/utils/SQLDataHandler.py index d60c806d2..f352ebea4 100755 --- a/src/utils/SQLDataHandler.py +++ b/src/utils/SQLDataHandler.py @@ -12,6 +12,7 @@ from config import config from config import global_vars +from MySQLDataHandler import record logger = logging.getLogger(__name__) @@ -20,7 +21,7 @@ sql_live_connect_num = 25 -class SQLConnManager: +class SQLConnManager(object): @staticmethod def Connect(): @@ -142,7 +143,7 @@ def ReturnConnection(conn): global_vars["sql_lock"].release() return None -class DataHandler: +class DataHandler(object): def __init__(self): start_time = timeit.default_timer() self.CreateDatabase() @@ -150,7 +151,7 @@ def __init__(self): logger.debug ("********************** created a new Data Handler *******************") self.conn = SQLConnManager.GetConnection() logger.debug ("Get database connection %s" % str(self.conn)) - + #print "Connecting to server ..." self.jobtablename = "jobs-%s" % config["clusterId"] self.acltablename = "acl-%s" % config["clusterId"] @@ -164,8 +165,6 @@ def __init__(self): elapsed = timeit.default_timer() - start_time logger.debug ("DataHandler initialization, time elapsed %f s" % elapsed) - - def CreateDatabase(self): if "initSQLDB" not in global_vars or not global_vars["initSQLDB"]: logger.info("===========init SQL database===============") @@ -327,7 +326,7 @@ def CreateTable(self): self.conn.commit() cursor.close() - + sql = """ if not exists (select * from sysobjects where name='%s' and xtype='U') CREATE TABLE [dbo].[%s] @@ -349,41 +348,34 @@ def CreateTable(self): self.conn.commit() cursor.close() - + @record def AddStorage(self, vcName, url, storageType, metadata, defaultMountPath): try: - start_time = timeit.default_timer() sql = "INSERT INTO [%s] (storageType, url, metadata, vcName, defaultMountPath) VALUES (?,?,?,?,?)""" % self.storagetablename cursor = self.conn.cursor() cursor.execute(sql, (storageType, url, metadata, vcName, defaultMountPath)) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: AddStorage to DB: url : %s, vcName: %s , time elapsed %f s" % (url, vcName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def DeleteStorage(self, vcName, url): try: - start_time = timeit.default_timer() sql = "DELETE FROM [%s] WHERE [url] = '%s' and [vcName] = '%s'" % (self.storagetablename, url, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: DeleteStorage: url:%s, vcName:%s, time elapsed %f s" % (url, vcName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def ListStorages(self, vcName): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [vcName],[url],[storageType],[metadata],[defaultMountPath] FROM [%s] WHERE [vcName] = '%s' " % (self.storagetablename, vcName) ret = [] @@ -402,45 +394,36 @@ def ListStorages(self, vcName): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: ListStorages time elapsed %f s" % (elapsed)) return ret - + @record def UpdateStorage(self, vcName, url, storageType, metadata, defaultMountPath): try: - start_time = timeit.default_timer() sql = """update [%s] set storageType = '%s', metadata = '%s', defaultMountPath = '%s' where [vcName] = '%s' and [url] = '%s' """ % (self.storagetablename, storageType, metadata, defaultMountPath, vcName, url) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: UpdateStorage: vcName: %s, url: %s, time elapsed %f s" % (vcName, url, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def AddVC(self, vcName, quota, metadata): try: - start_time = timeit.default_timer() sql = """INSERT INTO [%s] (vcName, quota, metadata) VALUES (?,?,?)""" % self.vctablename cursor = self.conn.cursor() cursor.execute(sql, vcName, quota, metadata) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: AddVC to DB: vcName: %s , time elapsed %f s" % (vcName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def ListVCs(self, vcName): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [vcName],[quota],[metadata] FROM [%s]" % (self.vctablename) ret = [] @@ -457,45 +440,36 @@ def ListVCs(self, vcName): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: ListVCs time elapsed %f s" % ( elapsed)) - return ret - + return ret + @record def DeleteVC(self, vcName): try: - start_time = timeit.default_timer() sql = "DELETE FROM [%s] WHERE [vcName] = '%s'" % (self.vctablename, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: DeleteVC: vcName: %s , time elapsed %f s" % (vcName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def UpdateVC(self, vcName, quota, metadata): try: - start_time = timeit.default_timer() sql = """update [%s] set quota = '%s', metadata = '%s' where [vcName] = '%s'""" % (self.vctablename, quota, metadata, vcName) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: UpdateVC: vcName: %s , time elapsed %f s" % (vcName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetIdentityInfo(self, identityName): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [identityName],[uid],[gid],[groups] FROM [%s] where [identityName] = '%s'" % (self.identitytablename, identityName) ret = [] @@ -513,35 +487,29 @@ def GetIdentityInfo(self, identityName): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: GetIdentityInfo time elapsed %f s" % (elapsed)) - return ret - + return ret + @record def UpdateIdentityInfo(self, identityName, uid, gid, groups): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() - + if len(self.GetIdentityInfo(identityName)) == 0: sql = """INSERT INTO [%s] (identityName,uid,gid,groups) VALUES (?,?,?,?)""" % self.identitytablename cursor.execute(sql, identityName, uid, gid, json.dumps(groups)) else: sql = """update [%s] set uid = '%s', gid = '%s', groups = '%s' where [identityName] = '%s' """ % (self.identitytablename, uid, gid, groups, identityName) cursor.execute(sql) - + self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: UpdateIdentityInfo %s to database , time elapsed %f s" % (identityName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetAceCount(self, identityId, resource): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM [%s] where [identityId] = '%s' and [resource] = '%s'" % (self.acltablename,identityId, resource) cursor.execute(query) @@ -550,88 +518,73 @@ def GetAceCount(self, identityId, resource): ret = c[0] self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: GetAceCount time elapsed %f s" % ( elapsed)) - return ret - + return ret + @record def UpdateAce(self, identityName, identityId, resource, permissions, isDeny): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() - + if self.GetAceCount(identityId, resource) == 0: sql = """INSERT INTO [%s] (identityName,identityId,resource,permissions,isDeny) VALUES (?,?,?,?,?)""" % self.acltablename cursor.execute(sql, identityName, identityId, resource, permissions, isDeny) else: sql = """update [%s] set permissions = '%s' where [identityName] = '%s' and [resource] = '%s' """ % (self.acltablename, permissions, identityName, resource) cursor.execute(sql) - + self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: UpdateAce %s - %s to database , time elapsed %f s" % (identityName, resource, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def UpdateAclIdentityId(self, identityName, identityId): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() sql = """update [%s] set identityName = '%s' where [identityName] = '%s' """ % (self.acltablename, identityId, identityName) cursor.execute(sql) - + self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: UpdateAclIdentityId %s - %s to database , time elapsed %f s" % (identityName, identityId, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def DeleteResourceAcl(self, resource): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() - + sql = "DELETE FROM [%s] WHERE [resource] = '%s'" % (self.acltablename, resource) cursor = self.conn.cursor() - + self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: DeleteResourceAcl %s, time elapsed %f s" % (resource, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def DeleteAce(self, identityName, resource): try: - start_time = timeit.default_timer() cursor = self.conn.cursor() - + sql = "DELETE FROM [%s] WHERE [identityName] = '%s' and [resource] = '%s'" % (self.acltablename, identityName, resource) cursor = self.conn.cursor() - + self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: DeleteAce %s : %s, time elapsed %f s" % (resource, identityName, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetAcl(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [identityName],[identityId],[resource],[permissions],[isDeny] FROM [%s]" % (self.acltablename) ret = [] @@ -650,13 +603,10 @@ def GetAcl(self): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: GetAcl time elapsed %f s" % ( elapsed)) - return ret - + return ret + @record def GetResourceAcl(self, resource): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [identityName],[identityId],[resource],[permissions],[isDeny] FROM [%s] where [resource] = '%s'" % (self.acltablename, resource) ret = [] @@ -675,30 +625,24 @@ def GetResourceAcl(self, resource): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: GetResourceAcl time elapsed %f s" % ( elapsed)) - return ret - + return ret + @record def AddJob(self, jobParams): try: - start_time = timeit.default_timer() sql = """INSERT INTO [%s] (jobId, familyToken, isParent, jobName, userName, vcName, jobType,jobParams ) VALUES (?,?,?,?,?,?,?)""" % self.jobtablename cursor = self.conn.cursor() jobParam = base64.b64encode(json.dumps(jobParams)) cursor.execute(sql, jobParams["jobId"], jobParams["familyToken"], jobParams["isParent"], jobParams["jobName"], jobParams["userName"], jobParams["vcName"], jobParams["jobType"],jobParam) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: added job %s to database, time elapsed %f s" % (jobParams["jobId"],elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or")): - start_time = timeit.default_timer() ret = [] cursor = self.conn.cursor() try: @@ -715,16 +659,14 @@ def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or" else: status_list = [ " [jobStatus] %s '%s' " % (op[0],s) for s in status.split(',')] status_statement = (" "+op[1]+" ").join(status_list) - query += " and ( %s ) " % status_statement + query += " and ( %s ) " % status_statement query += " order by [jobTime] Desc" - start_time1 = timeit.default_timer() cursor.execute(query) - elapsed1 = timeit.default_timer() - start_time1 - start_time2 = timeit.default_timer() + fetch_start = timeit.default_timer() data = cursor.fetchall() - elapsed2 = timeit.default_timer() - start_time2 - logger.info ("(fetchall time: %f)" % (elapsed2)) + fetch_time = timeit.default_timer() - fetch_start + logger.info ("(fetchall time: %f)" % (fetch_time)) for (jobId,jobName,userName, vcName,jobStatus,jobStatusDetail, jobType, jobDescriptionPath, jobDescription, jobTime, endpoints, jobParams,errorMsg, jobMeta) in data: record = {} record["jobId"] = jobId @@ -744,16 +686,12 @@ def GetJobList(self, userName, vcName, num = None, status = None, op = ("=","or" ret.append(record) except Exception as e: logger.error('Exception: '+ str(e)) - pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get job list for user %s , time elapsed %f s (SQL query time: %f)" % (userName, elapsed, elapsed1)) return ret - + @record def GetJob(self, **kwargs): - start_time = timeit.default_timer() valid_keys = ["jobId", "familyToken", "isParent", "jobName", "userName", "vcName", "jobStatus", "jobType", "jobTime"] if len(kwargs) != 1: return [] key, expected = kwargs.popitem() @@ -767,29 +705,23 @@ def GetJob(self, **kwargs): ret = [dict(zip(columns, row)) for row in cursor.fetchall()] self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get job details with query %s=%s , time elapsed %f s" % (key, expected, elapsed)) return ret - + @record def AddCommand(self,jobId,command): try: - start_time = timeit.default_timer() sql = """INSERT INTO [%s] (jobId, command) VALUES (?,?)""" % self.commandtablename cursor = self.conn.cursor() cursor.execute(sql, jobId, command) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: add command to database, jobId: %s , time elapsed %f s" % (jobId, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetPendingCommands(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [id], [jobId], [command] FROM [%s] WHERE [status] = 'pending' order by [time]" % (self.commandtablename) cursor.execute(query) @@ -802,29 +734,23 @@ def GetPendingCommands(self): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get pending command , time elapsed %f s" % (elapsed)) - return ret - + return ret + @record def FinishCommand(self,commandId): try: - start_time = timeit.default_timer() sql = """update [%s] set status = 'run' where [id] = '%s' """ % (self.commandtablename, commandId) cursor = self.conn.cursor() cursor.execute(sql) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: set command %s as finished , time elapsed %f s" % (commandId, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) return False - + @record def GetCommands(self, jobId): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [time], [command], [status], [output] FROM [%s] WHERE [jobId] = '%s' order by [time]" % (self.commandtablename, jobId) cursor.execute(query) @@ -838,13 +764,10 @@ def GetCommands(self, jobId): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get command list for job %s , time elapsed %f s" % (jobId, elapsed)) - return ret - + return ret + @record def GetPendingJobs(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [jobId],[jobName],[userName], [vcName], [jobStatus], [jobType], [jobDescriptionPath], [jobDescription], [jobTime], [endpoints], [jobParams],[errorMsg] ,[jobMeta] FROM [%s] where [jobStatus] <> 'error' and [jobStatus] <> 'failed' and [jobStatus] <> 'finished' and [jobStatus] <> 'killed' order by [jobTime] DESC" % (self.jobtablename) cursor.execute(query) @@ -867,45 +790,36 @@ def GetPendingJobs(self): ret.append(record) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get pending jobs , time elapsed %f s" % (elapsed)) - return ret - + return ret + @record def SetJobError(self,jobId,errorMsg): try: - start_time = timeit.default_timer() sql = """update [%s] set jobStatus = 'error', [errorMsg] = ? where [jobId] = '%s' """ % (self.jobtablename,jobId) cursor = self.conn.cursor() cursor.execute(sql,errorMsg) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: set job %s error status in database, time elapsed %f s" % (jobId, elapsed)) return True except Exception as e: logger.error('Exception: '+ str(e)) - return False - + return False + @record def UpdateJobTextField(self,jobId,field,value): try: - start_time = timeit.default_timer() sql = """update [%s] set [%s] = ? where [jobId] = '%s' """ % (self.jobtablename,field, jobId) cursor = self.conn.cursor() cursor.execute(sql,value) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: update job %s, field %s , time elapsed %f s" % (jobId, field, elapsed)) return True except Exception, e: logger.error('Exception: '+ str(e)) return False - + @record def GetJobTextField(self,jobId,field): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [jobId], [%s] FROM [%s] where [jobId] = '%s' " % (field, self.jobtablename,jobId) ret = None @@ -918,12 +832,10 @@ def GetJobTextField(self,jobId,field): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get filed %s of job %s , time elapsed %f s" % (field, jobId, elapsed)) return ret + @record def AddandGetJobRetries(self,jobId): - start_time = timeit.default_timer() sql = """update [%s] set [retries] = [retries] + 1 where [jobId] = '%s' """ % (self.jobtablename, jobId) cursor = self.conn.cursor() cursor.execute(sql) @@ -939,30 +851,24 @@ def AddandGetJobRetries(self,jobId): ret = value self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get and update retries for job %s , time elapsed %f s" % (jobId, elapsed)) return ret - + @record def UpdateClusterStatus(self,clusterStatus): try: - start_time = timeit.default_timer() sql = """INSERT INTO [%s] (status) VALUES (?)""" % self.clusterstatustablename cursor = self.conn.cursor() status = base64.b64encode(json.dumps(clusterStatus)) cursor.execute(sql,status) self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: update cluster status, time elapsed %f s" % (elapsed)) return True except Exception, e: logger.error('Exception: '+ str(e)) return False - + @record def GetClusterStatus(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT TOP 1 [time], [status] FROM [%s] order by [time] DESC" % (self.clusterstatustablename) ret = None @@ -977,13 +883,10 @@ def GetClusterStatus(self): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get cluster status , time elapsed %f s" % (elapsed)) return ret, time - + @record def GetUsers(self): - start_time = timeit.default_timer() cursor = self.conn.cursor() query = "SELECT [identityName],[uid] FROM [%s]" % (self.identitytablename) ret = [] @@ -996,11 +899,9 @@ def GetUsers(self): pass self.conn.commit() cursor.close() - elapsed = timeit.default_timer() - start_time - logger.info ("DataHandler: get users, time elapsed %f s" % ( elapsed)) return ret - + @record def GetActiveJobsCount(self): cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM [%s] where [jobStatus] <> 'error' and [jobStatus] <> 'failed' and [jobStatus] <> 'finished' and [jobStatus] <> 'killed' " % (self.jobtablename) @@ -1011,8 +912,9 @@ def GetActiveJobsCount(self): self.conn.commit() cursor.close() - return ret + return ret + @record def GetALLJobsCount(self): cursor = self.conn.cursor() query = "SELECT count(ALL id) as c FROM [%s]" % (self.jobtablename) @@ -1023,7 +925,7 @@ def GetALLJobsCount(self): self.conn.commit() cursor.close() - return ret + return ret def __del__(self): logger.debug("********************** deleted a DataHandler instance *******************") @@ -1039,7 +941,7 @@ def Close(self): CREATE_TABLE = False CREATE_DB = True dataHandler = DataHandler() - + if TEST_INSERT_JOB: jobParams = {} jobParams["id"] = "dist-tf-00001" @@ -1047,9 +949,8 @@ def Close(self): jobParams["user-id"] = "hongzl" jobParams["job-meta-path"] = "/dlws/jobfiles/***" jobParams["job-meta"] = "ADSCASDcAE!EDASCASDFD" - + dataHandler.AddJob(jobParams) - if CREATE_TABLE: dataHandler.CreateTable()