Lambda + awswrangler: Poor performance while handling "large" parquet files

Viewed 391

I'm currently writing a Lambda function to read parquet files from 100MB o 200MB on average using Python and the AWS wrangler function. The idea is to read the files and transform them to csv:

import awswrangler as wr
from io import StringIO

print('Loading function')

s3 = boto3.client('s3')
dest_bucket = "mydestbucket"

def lambda_handler(event, context):
    # Get the object from the event and show its content type
    bucket = event['Records'][0]['s3']['bucket']['name']
    key = urllib.parse.unquote_plus(event['Records'][0]['s3']['object']['key'], encoding='utf-8')
    try:
        response = s3.get_object(Bucket=bucket, Key=key)
        print("CONTENTO TYPE: " + response['ContentType'])

        if key.endswith('.parquet'):
            dfs = wr.s3.read_parquet(path=['s3://' + bucket + '/' + key], chunked=True, use_threads=True)
            
            count=0
            for df in dfs:
                csv_buffer = StringIO()
                df.to_csv(csv_buffer)
                s3_resource = boto3.resource('s3')
                #s3_resource.Object(dest_bucket, 'dfo.csv').put(Body=df)
                s3_resource.Object(dest_bucket, 'dfo_' + str(count) + '.csv').put(Body=csv_buffer.getvalue())
                count += 1
                
            return "File written" 

The function works ok when I use small files, but once I try with large files (100MB) it gives a timeout.

I already allocated 3GB of memory for Lambda and set a timeout of 10 min, however, it doesn't seem to do the trick.

Do you know how to improve the performance apart from allocating more memory?

Thanks!

2 Answers

I resolved the issue by creating a Layer using fastparquet which handles the memory in a more optimal fashion than aws wrangler

from io import StringIO
from datetime import datetime

import boto3
import fastparquet as fp
import s3fs
import urllib.parse


    #S3 fs initialization
    s3_fs = s3fs.S3FileSystem()
    fs = s3fs.core.S3FileSystem()

    s3fs_path = fs.glob(path=s3_path)
    my_open = s3_fs.open

    # Read parquet object using fastparquet
    fp_obj = fp.ParquetFile(s3fs_path, open_with=my_open)

    # Filter columns and build a pandas df
    new_df = fp_obj.to_pandas()

    # csv buffer to perform the parquet --> csv transformation
    csv_buffer = StringIO()
    new_df.to_csv(csv_buffer)

    s3_resource.Object(
            dest_bucket,
            f"{file_path}",
        ).put(Body=csv_buffer.getvalue())
Related