import boto3 import os import json import time import subprocess from botocore.config import Config as BotoCoreConfig from dotenv import load_dotenv # Load env file if it exists load_dotenv(verbose=True) ENVIRONMENT = os.environ.get('Environment', 'dev') region = os.environ.get('AWS_REGION') boto_config = BotoCoreConfig(read_timeout=70, region_name=region) sf_client = boto3.client('stepfunctions', config=boto_config) activity_arn = os.environ.get('SPLIT_OUTFILE_ARN') OUTFILE_DIR = os.environ.get('OUTFILE_DIR') SPLIT_LINES_BATCH_SIZE = os.environ.get('SPLIT_LINES_BATCH_SIZE') while True: response = sf_client.get_activity_task(activityArn=activity_arn, workerName=ENVIRONMENT+'_split_outfile') if 'taskToken' not in response: print('No Task Token') time.sleep(5) else: print(response['taskToken']) print("===================") activity_token = response['taskToken'] try: activity_input = response['input'].strip('\"') print(activity_input) completed_process = subprocess.run(['split', '-l', SPLIT_LINES_BATCH_SIZE, '--filter=/usr/bin/pigz > $FILE.gz', OUTFILE_DIR + activity_input + '/raw/outfile.txt', OUTFILE_DIR + activity_input + '/file_parts/part_decomp_'], check=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE) print(completed_process.stdout) except subprocess.CalledProcessError as exc: sf_client.send_task_failure(taskToken=activity_token, error=str(exc.stderr)) else: sf_client.send_task_success(taskToken=activity_token, output=json.dumps(activity_input))