如何按ma_id与items.nan分组PyArrow表并聚合最大processing_ts
问题:PyArrow中按列表嵌套字段分组求时间最大值
我已将Parquet文件加载到具有如下Schema的PyArrow表中:
import pyarrow as pa schema: pa.Schema = pa.schema( [ ("ma_id", pa.int32()), ("processing_ts", pa.timestamp("ms")), ( "items", pa.list_( pa.struct( [ pa.field("nan", pa.int32()), pa.field("ean", pa.int32()), ] ) ), ), ] )
表中示例数据如下:
[ (100, '2025-01-03 16:21:00', [{'nan': 1, 'ean': 11}, {'nan': 2, 'ean': 212}, {'nan': 3, 'ean': 3}]), (100, '2025-01-03 23:55:00', [{'nan': 9, 'ean': 95}, {'nan': 2, 'ean': 212}, {'nan': 9, 'ean': 95}]), (120, '2025-01-03 21:21:00', [{'nan': 8, 'ean': 87}, {'nan': 2, 'ean': 212}, {'nan': 9, 'ean': 95}]), (100, '2025-01-03 01:45:00', [{'nan': 6, 'ean': 666}, {'nan': 1, 'ean': 11}, {'nan': 7, 'ean': 711}, {'nan': 6, 'ean': 666}]), (120, '2025-01-03 12:38:00', [{'nan': 8, 'ean': 87}, {'nan': 9, 'ean': 95}]), ]
需求是按ma_id和items列表中的nan字段分组,获取每组的max(processing_ts),期望输出结果如下:
| ma_id | nan | max_processing_ts |
|---|---|---|
| 100 | 1 | '2025-01-03 16:21:00' |
| 100 | 2 | '2025-01-03 23:55:00' |
| 100 | 3 | '2025-01-03 16:21:00' |
| 100 | 6 | '2025-01-03 01:45:00' |
| 100 | 7 | '2025-01-03 01:45:00' |
| 100 | 9 | '2025-01-03 23:55:00' |
| 120 | 2 | '2025-01-03 21:21:00' |
| 120 | 8 | '2025-01-03 21:21:00' |
| 120 | 9 | '2025-01-03 21:21:00' |
解决方案
要实现嵌套列表字段的分组,核心是先展开列表字段,将每个items中的struct元素拆分为单独行,再进行分组聚合:
步骤1:展开列表字段
使用PyArrow的flatten方法展开items列表,将每个嵌套的struct转为独立行,同时保留原行的ma_id和processing_ts:
import pyarrow as pa import pyarrow.compute as pc # 假设已加载数据到table变量 # table = pa.parquet.read_table("your_file.parquet", schema=schema) # 展开items列表 flattened_table = table.flatten(["items"])
步骤2:提取nan字段并分组聚合
从展开后的struct字段中提取nan,然后按ma_id和nan分组,计算每组的最大processing_ts:
# 提取nan字段并重命名 table_with_nan = flattened_table.select( ["ma_id", "processing_ts", pc.field("items.nan").alias("nan")] ) # 分组聚合求最大时间 result = table_with_nan.group_by(["ma_id", "nan"]).aggregate( [("processing_ts", "max")] ).rename_columns(["ma_id", "nan", "max_processing_ts"])
验证结果
执行上述代码后,result表的结构和数据将与期望输出一致。可通过result.to_pandas()查看格式化后的表格结果。
内容的提问来源于stack exchange,提问作者Najib Bakahoui
相关产品推荐
相关产品推荐

