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

如何按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_idnanmax_processing_ts
1001'2025-01-03 16:21:00'
1002'2025-01-03 23:55:00'
1003'2025-01-03 16:21:00'
1006'2025-01-03 01:45:00'
1007'2025-01-03 01:45:00'
1009'2025-01-03 23:55:00'
1202'2025-01-03 21:21:00'
1208'2025-01-03 21:21:00'
1209'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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 19:50:08