"""Methods for processing CSV and XLSX via pandas dataframes. Allows for S3 access and S3->Snowflake data loading.""" import os from io import StringIO, BytesIO # python3; python2: BytesIO import pandas as pd from integration_scripts.connectors.logging import logger from integration_scripts import s3_backoff_utils as s3_utils from integration_scripts.connectors.s3 import s3_resource, s3_bucket_name def read_xlsx_to_dataframe_s3( source_key, bucket_name, override_na_values=None): """Load a XLSX s3 object into a dataframe. Parameters ---------- source_key (str) The key to load as a table bucket_name (str) The bucket containing the key override_na_values (list|None) A list of values with which to override the N/A value """ # Grab XLSX from S3 obj = s3_utils.get_object_backoff(bucket_name=bucket_name, key=source_key) # Load Excel as panda DataFrame if override_na_values is not None: df = pd.read_excel( obj['Body'], keep_default_na=False, na_values=override_na_values) else: df = pd.read_excel(obj['Body']) return df def convert_dataframe_to_csv_file(dframe, filename, drop_by=None): """Convert a DataFrame to an CSV file and save it to a file.""" logger.info('Writing DataFrame to CSV file: {}', filename) # Drop blank rows if drop_by: df2 = dframe.dropna(axis='rows', how='all', subset=[drop_by]) else: df2 = dframe.dropna(axis='rows', how='all') with open(filename, 'w', newline='') as fp: df2.to_csv(fp, index=False) logger.info('{} written.', filename) def convert_dataframe_to_csv_s3( dframe, output_path_and_key, drop_by=None): """Convert a DataFrame to an CSV file and save it to an S3 object.""" # Create target key target_key = s3_utils.get_s3_file_key(output_path_and_key + '.csv') logger.info('Writing DataFrame to CSV on S3: {}', target_key) # Drop blank rows if drop_by: df2 = dframe.dropna(axis='rows', how='all', subset=[drop_by]) else: df2 = dframe.dropna(axis='rows', how='all') # Create CSV format as a string buffer csv_buffer = StringIO() df2.to_csv(csv_buffer, index=False) # Save string buffer to S3 Bucket as csv file s3_resource.Object( s3_bucket_name, target_key).put(Body=csv_buffer.getvalue()) logger.info('{} written.', target_key) return target_key def convert_dataframe_to_xlsx_s3(dframe, output_path_and_file, drop_by=None): """Convert a DataFrame to an XLSX file and save it.""" output_path_and_file = os.path.splitext(output_path_and_file)[0] source_key = s3_utils.get_s3_file_key(output_path_and_file + '.csv') target_key = s3_utils.get_s3_file_key(output_path_and_file + '.xlsx') logger.info('Writing XLSX: {}', target_key) # Drop blank rows if drop_by: df2 = dframe.dropna(axis='rows', how='all', subset=[drop_by]) else: df2 = dframe.dropna(axis='rows', how='all') # Convert DataFrame to XLSX xlsx_buffer = BytesIO() df2.to_excel(xlsx_buffer, index=False) # Save XLSX stream to S3 bucket s3_resource.Object( s3_bucket_name, target_key).put(Body=xlsx_buffer.getvalue()) logger.info('{} saved.', target_key) return source_key, target_key def get_files_from_s3_bucket(folder_key, ext='csv', ignore_folders=None): """Create a dict of files, by directory, from S3.""" if not ignore_folders: ignore_folders = list() all_files = dict() # Filter to get s3 leaf keys folder_list = s3_utils.filter_backoff(folder_key, s3_bucket_name) for item in folder_list: file_path, extension = os.path.splitext(item.key) is_file_type = extension == '.' + ext is_directory = item.key.endswith('/') is_ignored = True if not is_directory: output_leaf_path = file_path.split('/')[-3] is_ignored = output_leaf_path in ignore_folders if not is_directory and is_file_type and not is_ignored: directory = item.key.split('/')[-2] if not all_files.get(directory): all_files[directory] = list() # file = item.key.split('/')[2] all_files[directory].append(item) return all_files def get_top_level_files_from_s3_bucket( folder_key, ext='csv', ignore_folders=None): """Create a dict of top-level files, by directory, from S3.""" if not ignore_folders: ignore_folders = list() all_files = list() # Filter to get s3 leaf keys folder_list = s3_utils.filter_backoff(folder_key, s3_bucket_name) for item in folder_list: file_path, extension = os.path.splitext(item.key) if os.path.split(file_path)[0] == folder_key: is_file_type = extension == '.' + ext is_directory = item.key.endswith('/') is_ignored = True if not is_directory: output_leaf_path = file_path.split('/')[-3] is_ignored = output_leaf_path in ignore_folders if not is_directory and is_file_type and not is_ignored: # directory = item.key.split('/')[-2] # all_files[directory] = list() # file = item.key.split('/')[2] all_files.append(item) return all_files def get_filename_key_from_s3_obj(obj, cut_extension=False): """Slice the filename (optionally w/o ext) from S3 obj.""" file_split = [r for r in reversed(obj.key.split('/'))] if cut_extension: # remove the extension filename_key = os.path.splitext(file_split[0])[0] else: filename_key = file_split[0] return filename_key def convert_all_csv_to_xlsx_s3(source_folder_key, output_folder_key=None, top_level=False, ignore_folders=None): """Convert all CSV to XLSX. """ if not output_folder_key: output_folder_key = source_folder_key # To return a list of new converted files converted_list = list() new_file_list = list() if top_level: all_files = dict() # Transform response object to dict with labels and files. all_files['#ROOT#'] = get_top_level_files_from_s3_bucket( source_folder_key, ext='csv') else: # Transform response object to dict with labels and files. all_files = get_files_from_s3_bucket( source_folder_key, ext='csv', ignore_folders=ignore_folders) # Loop through folders (label id's) for label, file_list in all_files.items(): # label_files = list() # Loop through files in each folder for file in file_list: # get the CSV object from bucket csv = s3_utils.get_object_backoff( bucket_name=s3_bucket_name, key=file.key) logger.info('Reading CSV to DataFrame: {}', file.key) # convert obj to pandas DataFrame dframe = pd.read_csv(csv['Body']) logger.info('CSV read.') # Get output file key output_file_key = get_filename_key_from_s3_obj( file, cut_extension=True) # Check if needed to add `label` folder to output_folder_key. if label == '#ROOT#': output_folder_key_label = s3_utils.get_s3_file_key( output_folder_key) else: output_folder_key_label = s3_utils.get_s3_file_key( output_folder_key, label) target_path = os.path.join(output_folder_key_label, output_file_key) try: # Drop Row Number column dframe.drop('ROW_NUMBER', 1, inplace=True) # Drop Completed Column dframe.drop('PROGRESS', 1, inplace=True) except KeyError: pass converted_file, target_key = convert_dataframe_to_xlsx_s3( dframe, target_path) converted_list.append(converted_file) new_file_list.append(target_key) # TODO: Figure out how to put multiple sheets in one workbook # load data frame to dict with label # label_files.append((output_file_key, df)) # with pd.ExcelWriter('{}.xlsx'.format(label)) as writer: # for file in label_files: # file[1].to_excel( # writer, sheet_name='{}'.format( # file[2].split('/')[3][:-4]), index=False) # print('test') # writer.save() return converted_list, new_file_list def convert_all_xlsx_to_csv_s3( source_folder_key, output_folder_key=None, sort_order=None, prune_list=None, drop_by=None): """Convert all XLSX to CSV.""" converted_file_list = list() if not output_folder_key: output_folder_key = source_folder_key # get all keys in folder files_only = s3_utils.filter_file_keys( source_folder_key, s3_bucket_name, file_ext='xlsx') # Loop through leaf keys for k in files_only: converted_file_list.append(k) # Create fully qualified keys source_key = s3_utils.get_s3_file_key(source_folder_key, k) # Read XLSX to dataframe logger.info('Converting file to DataFrame: {}.', source_key) dframe = read_xlsx_to_dataframe_s3(source_key, s3_bucket_name) logger.info('{} converted.', source_key) # TODO: Extract prune and sort to validation method. # Prune irrelevant fields. if prune_list: file_columns = [n for n in dframe.columns] drop_labels = list(set(file_columns) - set(prune_list)) if len(drop_labels): logger.info('Pruning irrelevant DataFrame columns.') dframe.drop(columns=drop_labels, inplace=True) logger.info('DataFrame pruned.') # Sort if valid sort order is passed. if type(sort_order) in [str, list]: sort_string = sort_order if type(sort_string) == list: sort_string = ', '.join(sort_string) logger.info('Sorting DataFrame by : {}.', sort_string) dframe.sort_values(by=sort_order, inplace=True) logger.info('DataFrame sorted.') output_file_key = os.path.splitext(os.path.split(source_key)[1])[0] output_path_and_key = s3_utils.get_s3_file_key(output_folder_key, output_file_key) convert_dataframe_to_csv_s3( dframe, output_path_and_key, drop_by=drop_by) return converted_file_list