Python低内存读取Parquet文件:逐行读取可行吗?适配S3与Lambda
解决方案:Lambda流式处理S3 Parquet文件并发送到Kafka
Parquet是列存储格式,不存在原生的"逐行读取"概念,但可以通过小批次读取行组内的数据实现近似流式的内存友好处理,同时兼容S3的StreamingBody。以下是针对需求的具体实现方案:
核心实现思路
利用PyArrow库的ParquetFile和scan_pandas接口,即使文件只有单个行组,也可以指定批次大小拆分读取,避免一次性加载全表到内存。同时直接将S3的StreamingBody转换为PyArrow可处理的流对象。
具体代码实现
import boto3 import pyarrow.parquet as pq from pyarrow import BufferReader from kafka import KafkaProducer import json def lambda_handler(event, context): # 初始化S3客户端 s3 = boto3.client('s3') bucket = event['bucket'] key = event['key'] # 获取S3文件的StreamingBody response = s3.get_object(Bucket=bucket, Key=key) stream_body = response['Body'] # 将StreamingBody转换为PyArrow BufferReader buffer = BufferReader(stream_body.read()) # 初始化ParquetFile对象 parquet_file = pq.ParquetFile(buffer) # 设置每批次读取的行数(根据Lambda内存配置调整,比如1000行/批) batch_size = 1000 # 初始化Kafka生产者(根据集群配置调整) producer = KafkaProducer( bootstrap_servers=['your-kafka-broker:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8') ) # 用scan_pandas直接生成批次迭代器,简化代码 scanner = parquet_file.scan_pandas(chunksize=batch_size) for batch_df in scanner: for _, record in batch_df.iterrows(): # 替换为实际的partition_key提取逻辑 partition_key = record['your_partition_key_column'] json_record = record.to_dict() # 构造目标格式消息并发送 message = f"{partition_key}: {json.dumps(json_record)}" producer.send('your-kafka-topic', value=message) # 确保所有消息发送完成 producer.flush() return {"statusCode": 200, "message": "Processing completed"}
关键细节说明
- StreamingBody适配:通过
BufferReader(stream_body.read())将S3流对象转换为PyArrow可读取的缓冲流,无需下载整个文件到本地。 - 内存控制:通过
chunksize参数控制每批次读取的行数,Lambda内存紧张时可进一步调小该值,避免OOM。 - 替代方案:行组拆分读取:如果需要更精细控制行组内的读取范围,也可以用
read_row_group结合Pandas切片实现:for row_group_idx in range(parquet_file.num_row_groups): row_group = parquet_file.read_row_group(row_group_idx) df = row_group.to_pandas() for start in range(0, len(df), batch_size): batch_df = df.iloc[start:start+batch_size] # 处理逻辑同上
之前尝试问题的原因
- fastparquet的
iter_row_groups是按行组粒度迭代,单个行组时自然会加载全量数据,无法拆分。 - PyArrow的
BufferReader没有readlines方法,是因为Parquet是列存格式,不能像CSV等行文本文件那样按行分割读取,必须通过批次读取适配列存特性。
内容的提问来源于stack exchange,提问作者James Kelleher
相关产品推荐
相关产品推荐

