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

如何用Apache Beam构造含唯一ID的新列并输出为JSON文件?

Apache Beam 实现分组生成唯一ID并输出换行JSON

完整代码示例

import apache_beam as beam
import json
import hashlib
from apache_beam.options.pipeline_options import PipelineOptions

def generate_group_key(row):
    # 替换成你用来分组的多列索引,比如取第0、1列作为分组依据
    # row是你已转换好的CSV字符串列表
    group_columns = (row[0], row[1])
    return (group_columns, row)

def generate_unique_id_and_format(grouped_data):
    group_key, rows = grouped_data
    # 方式1:用分组键生成哈希ID(低冲突,适合无意义唯一标识)
    key_str = "-".join(group_key).encode('utf-8')
    unique_id = hashlib.md5(key_str).hexdigest()
    
    # 方式2:直接用分组键作为ID(可读性强,适合键本身唯一的场景)
    # unique_id = "-".join(group_key)
    
    # 构造输出结构:唯一ID + 所属对象数组
    return {
        "id": unique_id,
        "items": rows
    }

def main():
    options = PipelineOptions()
    with beam.Pipeline(options=options) as p:
        # 1. 读取并转换CSV为字符串列表(你已实现的部分,这里用示例数据替代)
        # 实际场景替换为你的CSV读取逻辑,比如beam.io.ReadFromText后转列表
        csv_rows = p | "Create sample data" >> beam.Create([
            ["user1", "groupA", "data1"],
            ["user1", "groupA", "data2"],
            ["user2", "groupB", "data3"],
            ["user2", "groupB", "data4"]
        ])
        
        # 2. 生成分组键并执行分组
        grouped = (
            csv_rows
            | "Generate group key" >> beam.Map(generate_group_key)
            | "Group by key" >> beam.GroupByKey()
        )
        
        # 3. 生成唯一ID并构造JSON格式
        formatted_output = grouped | "Format output" >> beam.Map(generate_unique_id_and_format)
        
        # 4. 转换为JSON字符串并写入换行分隔文件
        (
            formatted_output
            | "Convert to JSON string" >> beam.Map(json.dumps)
            | "Write to newline-delimited JSON" >> beam.io.WriteToText(
                file_path_prefix="output/grouped_data",
                file_name_suffix=".json",
                num_shards=1,  # 生成单个文件,按需调整
                shard_name_template=""  # 移除分片后缀
            )
        )

if __name__ == "__main__":
    main()

关键步骤说明

  • 分组键生成:generate_group_key函数将你指定的多列转为可哈希的元组,作为GroupByKey的分组依据(Beam要求分组键必须是可哈希类型)。
  • 唯一ID生成:提供两种方案,哈希ID适合需要无意义唯一标识的场景,直接用分组键则更具可读性。
  • JSON输出:通过json.dumps将每个分组的字典转为字符串,WriteToText默认每行写入一个元素,正好满足换行分隔JSON的要求。

注意事项

  • 若你需要将唯一ID插入到原始行的索引0位置(而非输出分组数组),可调整generate_unique_id_and_format逻辑,给每个行添加ID后再输出,但当前代码更匹配你描述的“每行对应一个唯一ID及其所属对象数组”需求。
  • 处理大规模数据时,不建议强制num_shards=1,可根据集群资源调整分片数提升效率。
  • 哈希ID存在极小冲突概率,若需绝对唯一,可改用uuid.uuid4()生成UUID,注意每个分组仅生成一次即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 03:52:58