From 9fd0bbaaf9bfa11811b4344b22cf41c609e68892 Mon Sep 17 00:00:00 2001 From: Sheng Zha Date: Mon, 22 Jun 2020 11:14:18 -0700 Subject: [PATCH 1/2] AWS batch job tool for GluonNLP --- tools/batch/docker/Dockerfile | 27 +++++ tools/batch/docker/gluon_nlp_job.sh | 33 ++++++ tools/batch/submit-job.py | 154 ++++++++++++++++++++++++++++ tools/batch/wait-job.py | 93 +++++++++++++++++ 4 files changed, 307 insertions(+) create mode 100644 tools/batch/docker/Dockerfile create mode 100755 tools/batch/docker/gluon_nlp_job.sh create mode 100644 tools/batch/submit-job.py create mode 100644 tools/batch/wait-job.py diff --git a/tools/batch/docker/Dockerfile b/tools/batch/docker/Dockerfile new file mode 100644 index 0000000000..7122a7e013 --- /dev/null +++ b/tools/batch/docker/Dockerfile @@ -0,0 +1,27 @@ +FROM nvidia/cuda:10.1-cudnn7-devel-ubuntu18.04 + + RUN apt-get update && apt-get install -y --no-install-recommends \ + build-essential \ + locales \ + cmake \ + git \ + curl \ + vim \ + unzip \ + sudo \ + ca-certificates \ + libjpeg-dev \ + libpng-dev \ + libfreetype6-dev \ + libxft-dev &&\ + rm -rf /var/lib/apt/lists/* + + RUN curl -o ~/miniconda.sh -O https://repo.anaconda.com/miniconda/Miniconda3-latest-Linux-x86_64.sh && \ + chmod +x ~/miniconda.sh && \ + ~/miniconda.sh -b -p /opt/conda && \ + rm ~/miniconda.sh && \ + /opt/conda/bin/conda clean -ya + ENV PATH /opt/conda/bin:$PATH + RUN git clone https://github.com/dmlc/gluon-nlp + WORKDIR gluon-nlp + ADD gluon_nlp_job.sh . diff --git a/tools/batch/docker/gluon_nlp_job.sh b/tools/batch/docker/gluon_nlp_job.sh new file mode 100755 index 0000000000..9bd355b6b0 --- /dev/null +++ b/tools/batch/docker/gluon_nlp_job.sh @@ -0,0 +1,33 @@ +#!/bin/bash +date +echo "Args: $@" +env +echo "jobId: $AWS_BATCH_JOB_ID" +echo "jobQueue: $AWS_BATCH_JQ_NAME" +echo "computeEnvironment: $AWS_BATCH_CE_NAME" + +SOURCE_REF=$1 +CONDA_ENV=$2 +WORK_DIR=$3 +COMMAND=$4 +SAVED_OUTPUT=$5 +SAVE_PATH=$6 +REMOTE=$7 + +if [ ! -z $REMOTE ]; then + git remote set-url origin $REMOTE +fi; + +git fetch origin $SOURCE_REF:working +git checkout working +pip install -v -e .[extras] + +cd $WORK_DIR +/bin/bash -o pipefail -c "$COMMAND" +COMMAND_EXIT_CODE=$? +if [[ -f $SAVED_OUTPUT ]]; then + aws s3 cp $SAVED_OUTPUT s3://gluon-nlp-staging/$SAVE_PATH; +elif [[ -d $SAVED_OUTPUT ]]; then + aws s3 cp --recursive $SAVED_OUTPUT s3://gluon-nlp-staging/$SAVE_PATH; +fi; +exit $COMMAND_EXIT_CODE diff --git a/tools/batch/submit-job.py b/tools/batch/submit-job.py new file mode 100644 index 0000000000..ec99e44f47 --- /dev/null +++ b/tools/batch/submit-job.py @@ -0,0 +1,154 @@ +import argparse +import random +import re +import sys +import time +from datetime import datetime + +import boto3 +from botocore.compat import total_seconds + +parser = argparse.ArgumentParser(formatter_class=argparse.ArgumentDefaultsHelpFormatter) + +parser.add_argument('--profile', help='profile name of aws account.', type=str, + default=None) +parser.add_argument('--region', help='Default region when creating new connections', type=str, + default=None) +parser.add_argument('--name', help='name of the job', type=str, default='dummy') +parser.add_argument('--job-queue', help='name of the job queue to submit this job', type=str, + default='gluon-nlp-jobs') +parser.add_argument('--job-definition', help='name of the job job definition', type=str, + default='gluon-nlp-jobs:8') +parser.add_argument('--source-ref', + help='ref in GluonNLP main github. e.g. master, refs/pull/500/head', + type=str, default='master') +parser.add_argument('--work-dir', + help='working directory inside the repo. e.g. scripts/sentiment_analysis', + type=str, default='scripts/bert') +parser.add_argument('--saved-output', + help='output to be saved, relative to working directory. ' + 'it can be either a single file or a directory', + type=str, default='.') +parser.add_argument('--save-path', + help='s3 path where files are saved.', + type=str, default='batch/temp/{}'.format(datetime.now().isoformat())) +parser.add_argument('--conda-env', + help='conda environment preset to use.', + type=str, default='gpu/py3') +parser.add_argument('--command', help='command to run', type=str, + default='git rev-parse HEAD | tee stdout.log') +parser.add_argument('--remote', + help='git repo address. https://github.com/dmlc/gluon-nlp', + type=str, default="https://github.com/dmlc/gluon-nlp") +parser.add_argument('--wait', help='block wait until the job completes. ' + 'Non-zero exit code if job fails.', action='store_true') +parser.add_argument('--timeout', help='job timeout in seconds', default=None, type=int) + +args = parser.parse_args() + +session = boto3.Session(profile_name=args.profile, region_name=args.region) +batch, cloudwatch = [session.client(service_name=sn) for sn in ['batch', 'logs']] + +def printLogs(logGroupName, logStreamName, startTime): + kwargs = {'logGroupName': logGroupName, + 'logStreamName': logStreamName, + 'startTime': startTime, + 'startFromHead': True} + + lastTimestamp = 0 + while True: + logEvents = cloudwatch.get_log_events(**kwargs) + + for event in logEvents['events']: + lastTimestamp = event['timestamp'] + timestamp = datetime.utcfromtimestamp(lastTimestamp / 1000.0).isoformat() + print('[{}] {}'.format((timestamp + '.000')[:23] + 'Z', event['message'])) + + nextToken = logEvents['nextForwardToken'] + if nextToken and kwargs.get('nextToken') != nextToken: + kwargs['nextToken'] = nextToken + else: + break + return lastTimestamp + + +def getLogStream(logGroupName, jobName, jobId): + response = cloudwatch.describe_log_streams( + logGroupName=logGroupName, + logStreamNamePrefix=jobName + '/' + jobId + ) + logStreams = response['logStreams'] + if not logStreams: + return '' + else: + return logStreams[0]['logStreamName'] + +def nowInMillis(): + endTime = long(total_seconds(datetime.utcnow() - datetime(1970, 1, 1))) * 1000 + return endTime + + +def main(): + spin = ['-', '/', '|', '\\', '-', '/', '|', '\\'] + logGroupName = '/aws/batch/job' + + jobName = re.sub('[^A-Za-z0-9_\-]', '', args.name)[:128] # Enforce AWS Batch jobName rules + jobQueue = args.job_queue + jobDefinition = args.job_definition + command = args.command.split() + wait = args.wait + + parameters={ + 'SOURCE_REF': args.source_ref, + 'WORK_DIR': args.work_dir, + 'SAVED_OUTPUT': args.saved_output, + 'SAVE_PATH': args.save_path, + 'CONDA_ENV': args.conda_env, + 'COMMAND': args.command, + 'REMOTE': args.remote + } + kwargs = dict( + jobName=jobName, + jobQueue=jobQueue, + jobDefinition=jobDefinition, + parameters=parameters, + ) + if args.timeout is not None: + kwargs['timeout'] = {'attemptDurationSeconds': args.timeout} + submitJobResponse = batch.submit_job(**kwargs) + + jobId = submitJobResponse['jobId'] + print('Submitted job [{} - {}] to the job queue [{}]'.format(jobName, jobId, jobQueue)) + + spinner = 0 + running = False + status_set = set() + startTime = 0 + + while wait: + time.sleep(random.randint(5, 10)) + describeJobsResponse = batch.describe_jobs(jobs=[jobId]) + status = describeJobsResponse['jobs'][0]['status'] + if status == 'SUCCEEDED' or status == 'FAILED': + print('=' * 80) + print('Job [{} - {}] {}'.format(jobName, jobId, status)) + + sys.exit(status == 'FAILED') + + elif status == 'RUNNING': + logStreamName = getLogStream(logGroupName, jobName, jobId) + if not running: + running = True + print('\rJob [{} - {}] is RUNNING.'.format(jobName, jobId)) + if logStreamName: + print('Output [{}]:\n {}'.format(logStreamName, '=' * 80)) + if logStreamName: + startTime = printLogs(logGroupName, logStreamName, startTime) + 1 + elif status not in status_set: + status_set.add(status) + print('\rJob [%s - %s] is %-9s... %s' % (jobName, jobId, status, spin[spinner % len(spin)]),) + sys.stdout.flush() + spinner += 1 + +if __name__ == '__main__': + main() diff --git a/tools/batch/wait-job.py b/tools/batch/wait-job.py new file mode 100644 index 0000000000..87d8679255 --- /dev/null +++ b/tools/batch/wait-job.py @@ -0,0 +1,93 @@ +import argparse +from datetime import datetime +import sys +import time + +import boto3 +from botocore.compat import total_seconds + +parser = argparse.ArgumentParser(formatter_class=argparse.ArgumentDefaultsHelpFormatter) + +parser.add_argument('--profile', help='profile name of aws account.', type=str, + default=None) +parser.add_argument('--job-id', help='job id to check status and wait.', type=str, + default=None) + +args = parser.parse_args() + +session = boto3.Session(profile_name=args.profile) +batch, cloudwatch = [session.client(service_name=sn) for sn in ['batch', 'logs']] + +def printLogs(logGroupName, logStreamName, startTime): + kwargs = {'logGroupName': logGroupName, + 'logStreamName': logStreamName, + 'startTime': startTime, + 'startFromHead': True} + + lastTimestamp = 0 + while True: + logEvents = cloudwatch.get_log_events(**kwargs) + + for event in logEvents['events']: + lastTimestamp = event['timestamp'] + timestamp = datetime.utcfromtimestamp(lastTimestamp / 1000.0).isoformat() + print('[{}] {}'.format((timestamp + '.000')[:23] + 'Z', event['message'])) + + nextToken = logEvents['nextForwardToken'] + if nextToken and kwargs.get('nextToken') != nextToken: + kwargs['nextToken'] = nextToken + else: + break + return lastTimestamp + + +def getLogStream(logGroupName, jobName, jobId): + response = cloudwatch.describe_log_streams( + logGroupName=logGroupName, + logStreamNamePrefix=jobName + '/' + jobId + ) + logStreams = response['logStreams'] + if not logStreams: + return '' + else: + return logStreams[0]['logStreamName'] + +def nowInMillis(): + endTime = long(total_seconds(datetime.utcnow() - datetime(1970, 1, 1))) * 1000 + return endTime + + +def main(): + spin = ['-', '/', '|', '\\', '-', '/', '|', '\\'] + logGroupName = '/aws/batch/job' + + jobId = args.job_id + + spinner = 0 + running = False + startTime = 0 + + while True: + time.sleep(1) + describeJobsResponse = batch.describe_jobs(jobs=[jobId]) + job = describeJobsResponse['jobs'][0] + status, jobName = job['status'], job['jobName'] + if status == 'SUCCEEDED' or status == 'FAILED': + print('=' * 80) + print('Job [{} - {}] {}'.format(jobName, jobId, status)) + break + elif status == 'RUNNING': + logStreamName = getLogStream(logGroupName, jobName, jobId) + if not running and logStreamName: + running = True + print('\rJob [{} - {}] is RUNNING.'.format(jobName, jobId)) + print('Output [{}]:\n {}'.format(logStreamName, '=' * 80)) + if logStreamName: + startTime = printLogs(logGroupName, logStreamName, startTime) + 1 + else: + print('\rJob [%s - %s] is %-9s... %s' % (jobName, jobId, status, spin[spinner % len(spin)]),) + sys.stdout.flush() + spinner += 1 + +if __name__ == '__main__': + main() From 8adfdaae8e4ce63c955405a0bd7c3a8ee04f7318 Mon Sep 17 00:00:00 2001 From: Sheng Zha Date: Sat, 27 Jun 2020 17:20:21 -0700 Subject: [PATCH 2/2] limit range --- .github/workflows/unittests.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/unittests.yml b/.github/workflows/unittests.yml index c9f41becd1..097f80243f 100644 --- a/.github/workflows/unittests.yml +++ b/.github/workflows/unittests.yml @@ -35,7 +35,7 @@ jobs: python -m pip install --user --upgrade pip python -m pip install --user setuptools pytest pytest-cov python -m pip install --upgrade cython - python -m pip install --pre --user mxnet>=2.0.0b20200604 -f https://dist.mxnet.io/python + python -m pip install --pre --user 'mxnet>=2.0.0b20200604,<2.0.0b20200621' -f https://dist.mxnet.io/python python -m pip install --user -e .[extras] - name: Test project run: |