import pandas as pd import boto3 import sys def merge_tables_from_s3(folder, prefix, run_id): ''' Merge all csv files into a single one folder -- where work is stored prefix -- assigned prefix for outcome of merge run_id -- ID of current run ''' # Listing all tables for the given run_id table_names = [] for obj in bucket.objects.filter(Prefix=folder): table_name = obj.key.split('/')[-1] prefix_run_id = "{}_run_{}".format(prefix, run_id) if table_name.startswith(prefix_run_id): table_names.append(table_name) # Merging all of them in a single dataframe merged_df = pd.DataFrame() for table_name in table_names: obj = bucket.Object(f'{folder}/{table_name}') table_data = pd.read_csv(obj.get()['Body']) merged_df = pd.concat([merged_df,table_data]) print(merged_df.shape) return merged_df def save_dataframe_s3(df, folder, prefix, run_id): ''' Save the given parallel chunk result into S3: df -- dataframe to save folder -- where work is stored prefix -- assigned prefix for file name run_id -- ID of current run ''' s3 = boto3.client('s3') bucket_name = 'dev-cucumbers' filepath = "{}/{}_{}_merged.csv".format(folder, prefix, run_id) csv_buffer = df.to_csv(index=False).encode('utf-8') # Save the CSV file to S3 s3.put_object(Body=csv_buffer, Bucket=bucket_name, Key=filepath) print(f"Table saved to S3 bucket: {bucket_name}, with file name: {filepath}") def get_argument(args, name, default_value = None): ''' Getting arguments from processing script ''' arg_name = "--" + name if arg_name in args: index = args.index("--" + name) + 1 return args[index] else: return default_value def get_mandatory_argument(args, name): ''' Getting mandatory arguments from processing script ''' res = get_argument(args, name) if res is None: raise "Missing script mandatory argument --{}".format(name) else: return res if __name__ == '__main__': # Run id RUN_ID = get_mandatory_argument(sys.argv, "run-id") # Folder FOLDER = get_mandatory_argument(sys.argv, "folder") PREFIX_arima = 'arima_cross' PREFIX_fourier = 'fourier_table' PREFIX_ts = 'arima_timeseries' PREFIX_fourier_gb = 'series_with_inflection_last_week' s3 = boto3.resource('s3') bucket_name = 'dev-cucumbers' bucket = s3.Bucket(bucket_name) merged_df_arima = merge_tables_from_s3(FOLDER, PREFIX_arima, RUN_ID) merged_df_fourier = merge_tables_from_s3(FOLDER, PREFIX_fourier, RUN_ID) merged_df_ts = merge_tables_from_s3(FOLDER, PREFIX_ts, RUN_ID) merged_df_fourier_gb = merge_tables_from_s3(FOLDER, PREFIX_fourier_gb, RUN_ID) save_dataframe_s3(merged_df_arima, FOLDER, PREFIX_arima, RUN_ID) save_dataframe_s3(merged_df_fourier, FOLDER, PREFIX_fourier, RUN_ID) save_dataframe_s3(merged_df_ts, FOLDER, PREFIX_ts, RUN_ID) save_dataframe_s3(merged_df_fourier_gb, FOLDER, PREFIX_fourier_gb, RUN_ID)