import concurrent.futures import csv import boto3 import smart_open from ssavva.sync_to_dynamodb import config from ssavva.sync_to_dynamodb import utils dynamodb = boto3.resource('dynamodb') table = dynamodb.Table(config.dynamo_table_name) def get_products(s3_path): with smart_open.smart_open(s3_path, 'r') as f: reader = csv.reader(f, quotechar='"', delimiter=',') for r in reader: yield r def put_in_dynamo(items): with table.batch_writer() as batch: for item in items: batch.put_item(Item=dict(zip(config.product_fields, item))) def main(data_type): bucket = config.s3_bucket path = 'dynamo_sync/{}_full/'.format(data_type) s3_keys = utils.keys_in_s3_location(bucket, path) data_generators = ( (i, get_products(s3_path)) for i, s3_path in enumerate(s3_keys) ) with concurrent.futures.ThreadPoolExecutor() as executor: futures = [ executor.submit(put_in_dynamo, data, i) for i, data in data_generators ] for future in concurrent.futures.as_completed(futures): future.result() if __name__ == '__main__': main('products')