I have a function below, that iterates through an aws s3 bucket to capture values in each file to pass in the body for an api request. There are many api request per file and many files - so the function is as below? How can I best leverage this function alongside the concurrent library in python, to make request concurrently, while appropriately iterating through new values in file and new files within s3?
def clean_adresses(url,unload_bucket,cleaned_bucket):
pages = sgv.paginator.paginate(Bucket=glob_common_vars.s3_bucket_name, Prefix=unload_bucket)
final_output = []
for page in pages:
for obj in page['Contents']:
addresses = glob_common_vars.s3_resource.Object(glob_common_vars.s3_bucket_name,obj['Key'])
with gzip.GzipFile(fileobj=addresses.get()["Body"]) as gzipfile:
addresses = ndjson.loads(gzipfile.read().decode('utf-8'))
addresses_len = len(addresses)
b = 0
batch_num = 0
while b < addresses_len:
a = b
b = b + 100
final_addresses = addresses[a:b]
response = requests.request("POST", url, json=final_addresses)
resp_data = response.json()
print(resp_data)
final_output.extend(resp_data)
if len(final_output) > 10000:
gec.ndjson_to_s3(final_output, glob_common_vars.s3_resource, glob_common_vars.s3_bucket_name, cleaned_bucket + str(batch_num) + '.json.gzip')
batch_num +=1
final_output = []
print(final_output)
gec.ndjson_to_s3(final_output, glob_common_vars.s3_resource, glob_common_vars.s3_bucket_name, cleaned_bucket + str(batch_num) + '.json.gzip')
Here is how I am starting , but wondering how I can handle iterating through "addresses" ie the [a:b] and the new files, without making redundant request.
with concurrent.futures.ThreadPoolExecutor(max_workers = 5) as executor:
future_cleaned_addresses = {executor.submit(clean_adresses,.....?