Unable to make aiohttp requests on AWS Lambda due to "TypeError: A Future or coroutine is required"

Viewed 1200

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))
0 Answers
Related