You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

AWS Lambda分块读取S3大Parquet文件遇Schema不一致错误求助

解决AWS Lambda读取Parquet分块时的Schema不匹配错误

你遇到的ArrowInvalid: Schema at index 0 was different错误,本质是你读取的Parquet资源(单个大文件的不同row group,或者前缀下的多个Parquet文件)存在Schema不一致的情况——比如部分块的字段类型、字段数量和其他块不匹配,这种差异会在awswrangler分块读取时触发Schema校验失败,即使你尝试转字符串也没用,因为校验在读取阶段就已经执行。

1. 定位Schema差异根源

先运行以下代码,查看所有Parquet块的Schema,找出具体差异的字段:

import awswrangler as wr
# 替换为你的bucket和prefix
schemas = wr.s3.read_parquet_metadata(path=f"s3://{bucket_name}/{object_prefix}")
for idx, schema in enumerate(schemas):
    print(f"Schema {idx}: {schema}")

对比输出的Schema,就能看到哪些字段在不同块中类型不一致,或者存在缺失/新增的情况。

2. 强制统一Schema读取(推荐方案)

手动定义一个符合你需求的统一Schema,让awswrangler按这个Schema读取所有块,自动转换不匹配的字段:

步骤1:定义统一Schema

用pyarrow定义你需要的字段和对应类型(比如统一转成字符串,或者指定正确的类型):

import pyarrow as pa
# 根据你的columns_to_keep定义统一Schema,示例如下
unified_schema = pa.schema([
    ("nom_archive", pa.string()),
    ("nom_fichier", pa.string()),
    ("publication_date", pa.string()),
    # 其他需要保留的字段,按需添加并指定统一类型
])

步骤2:修改read_parquet调用

在读取时指定schema参数,强制使用统一Schema:

dfs = wr.s3.read_parquet(
    path=f"s3://{bucket_name}/{object_prefix}",
    columns=columns_to_keep,
    chunked=True,
    schema=unified_schema,  # 强制统一Schema
    coerce_int96_timestamp_unit="ns"  # 若存在timestamp类型不一致,添加此参数处理
)

3. 临时方案:关闭Schema校验(不推荐)

如果只是临时处理,且能接受潜在的数据异常,可以关闭awswrangler的Schema校验:

dfs = wr.s3.read_parquet(
    path=f"s3://{bucket_name}/{object_prefix}",
    columns=columns_to_keep,
    chunked=True,
    validate_schema=False  # 关闭Schema校验
)

额外优化建议

  • 你的逻辑在每个chunk最终只保留1行数据,可以考虑在读取阶段就过滤,减少内存占用(比如先获取nom_archive最大值对应的行,避免读取整个chunk)。
  • 替换df.iterrows()为df.to_dict('records'),提升数据转换效率:
    # 替换原有的循环逻辑
    processed_rows = df.to_dict('records')
    for row in processed_rows:
        rows.append({
            "grid_operator": "grid_operator",
            "file_name": row["nom_fichier"],
            "file_size": None,
            "archive_name": row["nom_archive"],
            "parent_path": "parent_path",
            "event_time": "event_time",
            "upload_day": "upload_day",
            "upload_month": "upload_month",
            "extraction_date": row.get("publication_date"),
            "flow_name": flow_name
        })
    

内容的提问来源于stack exchange,提问作者user24123007

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 06:44:52