I want to make async API get requests on my aws lambda (runtime is python 3.6). I originally tried to use grequests but I don't know how to compile that into a linux compatible bundle to import to AWS lambda (I would encounter the "Gevent is required for grequests error"). Instead, now I am using aysncio/aiohttp. It works completely fine on my local machine, but when I transfer it to AWS lambda (manually deploy by zipping packages and src folder to a .zip then uploading), I encounter
TypeError: A Future or coroutine is required
/var/runtime/awslambda/bootstrap.py:290: RuntimeWarning: coroutine 'main' was never awaited
Here are the relevant parts of my code:
import os
import boto3
import json
import base64
from datetime import datetime
import logging
import asyncio
import aiohttp
async def fetch(session, headers, url):
async with session.get(url, headers=headers ) as response:
logging.info(response.status)
return await response.text()
async def main(urls, headers):
tasks = []
connector = aiohttp.TCPConnector(limit=13)
async with aiohttp.ClientSession(connector=connector) as session:
for url in urls:
task = asyncio.ensure_future(fetch(session, headers, url))
tasks.append(task)
responses = await asyncio.gather(*tasks)
return responses
# Invoked method
def lambda_handler(event, context):
my_urls = []
for record in event['Records']:
decoded = base64.b64decode(record['kinesis']['data'])
try:
url_record = json.loads(decoded)
url = url_record['url']
patient_urls.append(url)
except (json.JSONDecodeError, KeyError):
error = "Could not read url from record: {}".format(decoded)
logging.error(error)
continue
token = get_token() # method not included as not relevant to error
if token is not None:
headers = {"Authorization": "Bearer {}".format(token)}
loop = asyncio.get_event_loop()
results = loop.run_until_complete(main(my_urls, headers))
print("we did it:",len(results))
else:
timestamp = str(datetime.now())
error = timestamp + "Token error"
logging.error(error)
I don't get why this error occurs as it works perfectly fine on my local machine and I practically copy-pasted the code. Here is the local code for reference:
async def fetch(session, headers, url):
async with session.get(url, headers=headers ) as response:
return await response.text()
async def main(urls, headers):
tasks = []
connector = aiohttp.TCPConnector(limit=13)
async with aiohttp.ClientSession(connector=connector) as session:
for url in urls:
task = asyncio.ensure_future(fetch(session, headers, url))
tasks.append(task)
responses = await asyncio.gather(*tasks)
return responses
start = time.time()
token = get_token()
headers = {"Authorization": "Bearer {}".format(token)}
patient_urls = get_urls()
loop = asyncio.get_event_loop()
a = loop.run_until_complete(main(patient_urls, headers))
print("we did it!",len(a))
end = time.time()
print("it took {} seconds".format(end-start))