Created
February 9, 2016 00:27
-
-
Save rodrickbrown/77bc9ac74ec6b7d43ee0 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env python3 | |
| import configparser | |
| import random | |
| import os | |
| import subprocess | |
| import argparse | |
| import sys | |
| import shlex | |
| import time | |
| from subprocess import check_call | |
| #from ConfigParser import * | |
| def readCfgOptions(section): | |
| cfgOptions = dict() | |
| for param in section: | |
| try: | |
| cfgOptions[param] = section.get(param) | |
| if cfgOptions[param] == -1: | |
| ebugPrint("skip: %s" % param) | |
| except: | |
| print("exception on %s!" % param) | |
| cfgOptions[param] = None | |
| return(cfgOptions) | |
| def buildCommandString(params, executor): | |
| execString = [] | |
| if executor == 'spark': | |
| execString.append('/usr/bin/timeout {}'.format(params.get('timeout'))) | |
| execString.append(params.get('sparkexec')) | |
| execString.append('--master {}'.format(params.get('mesosmasters'))) | |
| execString.append('--conf spark.ui.port={}'.format(os.environ.get('PORT0',"312" + str(random.randint(1,99))))) | |
| execString.append('--conf spark.mesos.coarse=true') | |
| execString.append('--conf spark.mesos.extra.cores={}'.format(params.get('cores'))) | |
| execString.append('--conf spark.cores.max={}'.format(params.get('coresmax'))) | |
| execString.append('--conf spark.mesos.constraints="rack:spark"') | |
| execString.append('--conf spark.sql.tungsten.enabled={}'.format(params.get('sparksqltungsten','true'))) | |
| execString.append('--class {}'.format(params.get('classname'))) | |
| execString.append('--total-executor-cores {}'.format(params.get('cores'))) | |
| execString.append('--driver-memory {}'.format(params.get('memory'))) | |
| execString.append('--executor-memory {}'.format(params.get('memory'))) | |
| execString.append('--jars {} {}'.format(params.get('classpath'), params.get('jarfile'))) | |
| if executor == 'java': | |
| execString.append('/usr/bin/timeout {}'.format(params.get('timeout'))) | |
| execString.append(params.get('javaexec')) | |
| execString.append('-Xmx{}'.format(params.get('memory'))) | |
| execString.append('-XX:+UseConcMarkSweepGC') | |
| execString.append('-XX:+CMSClassUnloadingEnabled') | |
| execString.append('-XX:+CMSParallelRemarkEnabled') | |
| execString.append('-XX:+PrintCommandLineFlags') | |
| execString.append('-Xss16M') | |
| execString.append('-Duser.timezone=GMT') | |
| execString.append('-cp {}:{}'.format(params.get('jarfile'), params.get('classpath'))) | |
| execString.append(params.get('classname')) | |
| return execString | |
| def executeJob(commandString, jarFile, job): | |
| print("\n*********************** Starting JOB: {} using SHA: {} ***********************\n" \ | |
| .format(job, os.path.basename('jarFile'))) | |
| print("\n{}\n\n".format(" ".join(commandString))) | |
| time.sleep(3) | |
| check_call(shlex.split(" ".join(commandString))) | |
| if __name__ == '__main__': | |
| parser = argparse.ArgumentParser(description="Orchard Batch Job Executor v1") | |
| parser.add_argument("-j", "--jobname", help="Name of Chronos Job to run", required=False) | |
| parser.add_argument("-l", "--listjobs", action="store_true", help="List all available jobs", required=False) | |
| args = parser.parse_args() | |
| cfgFile = "/tmp/load.cfg" | |
| config = configparser.ConfigParser(interpolation=configparser.ExtendedInterpolation()) | |
| commandString = {} | |
| if not args.jobname and not args.listjobs: | |
| parser.print_help() | |
| sys.exit(-1) | |
| if not os.access(cfgFile, os.R_OK): | |
| sys.stderr.write("Fatal error: {} is not readable\n".format(cfgFile)) | |
| sys.exit(-1) | |
| config.read(cfgFile) | |
| jobs = config.sections() | |
| if args.listjobs: | |
| jobs.remove('library') | |
| jobs.remove('spark') | |
| jobs.remove('java') | |
| for job in jobs: print(job) | |
| print("\n--------------\n{} jobs found".format(len(jobs))) | |
| sys.exit(0) | |
| if not args.jobname in jobs: | |
| sys.stderr.write("Fatal error: job \"{}\" not found in {}\n".format(args.jobname, cfgFile)) | |
| sys.exit(-1) | |
| print("Found {} jobs definitions in {}".format(len(jobs)-2, cfgFile)) | |
| jobCfg = {} | |
| sparkDefaults = {} | |
| javaDefaults = {} | |
| for job in jobs: | |
| if job == 'spark': | |
| sparkDefaults = readCfgOptions(config[job]) | |
| continue | |
| if job == 'java': | |
| javaDefaults = readCfgOptions(config[job]) | |
| continue | |
| if args.jobname == job: | |
| jobCfg = readCfgOptions(config[job]) | |
| if jobCfg.get('executor') == 'spark': | |
| jobCfg.update(sparkDefaults) | |
| commandString = buildCommandString(jobCfg, 'spark') | |
| if jobCfg.get('executor') == 'java': | |
| jobCfg.update(javaDefaults) | |
| commandString = buildCommandString(jobCfg, 'java') | |
| print("Loading {} from config: {}".format(args.jobname, cfgFile)) | |
| executeJob(commandString, jobCfg.get('jarfile'), job) | |
| ---- sample config file --- | |
| [spark] | |
| mesosMasters="mesos://zk://prod-zookeeper-1.aws.orchardplatform.com:2181,prod-zookeeper-2.aws.orchardplatform.com:2181,prod-zookeeper-3.aws.orchardplatform.com:2181/mesos" | |
| sparkExec=/opt/spark-1.5.1/bin/spark-submit | |
| memory=3G | |
| cores=3 | |
| coresMax=5 | |
| timeout=3600 | |
| [java] | |
| javaExec=/usr/bin/java | |
| memory=2048M | |
| classpath=/data/orchard/etc/raw-data-job-library | |
| jarFile=/data/orchard/jars/raw-data-job-library-3fd2628b53e812aec68adf1850317a3f8d1864d2-assembled.jar | |
| timeout=3600 | |
| ; (Warning) Please note modifying parms in the [java] or [spark] blocks will affect all jobs below. | |
| ; Jobs will inherit the default values based on the executor set java/spark | |
| ; values can also override default values if the options are set in their configuration block. | |
| [LoadTrade_Prosper] | |
| executor=spark | |
| className=com.orchard.dataloader.library.originators.prosper.LoadTrade_Prosper | |
| classPath=/data/orchard/etc/load-tradedata-accumulo-prod | |
| jarFile=/data/orchard/jars/dataloader-library-493e97b6c86b4b0adc7fd83514ce9e9d65b87278-assembled.jar | |
| [LoadAccountDetail_Prosper] | |
| executor=spark | |
| className=com.orchard.dataloader.library.originators.prosper.LoadAccountDetail_Prosper | |
| classPath=/data/orchard/etc/load-tradedata-accumulo-prod | |
| jarFile=/data/orchard/jars/dataloader-library-493e97b6c86b4b0adc7fd83514ce9e9d65b87278-assembled.jar | |
| [FtpToHdfsSync] | |
| executor=java | |
| className=com.orchard.rawdatajob.library.RawDataJob FtpToHdfsSync | |
| classPath=/data/orchard/etc/raw-data-job-library | |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment