Skip to content

Instantly share code, notes, and snippets.

@rodrickbrown
Created February 9, 2016 00:27
Show Gist options
  • Select an option

  • Save rodrickbrown/77bc9ac74ec6b7d43ee0 to your computer and use it in GitHub Desktop.

Select an option

Save rodrickbrown/77bc9ac74ec6b7d43ee0 to your computer and use it in GitHub Desktop.
#!/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