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

如何在Apache Beam中合并两个文件并查看生成的PCollection

问题排查

你的代码存在4个核心问题:

  • 变量名不匹配:上下文管理器中定义的管道实例名为pipeline,但后续读取数据时误用了未定义的p变量
  • 输入格式不符合CoGroupByKey要求:该转换要求所有输入的PCollection都为(键, 值)的二元组格式,直接读取的文本行是字符串,没有提取关联键user_id
  • 未跳过CSV表头:表头行也会被当作普通数据处理,会导致解析异常
  • 示例数据无匹配关联键:你提供的两个CSV样例中,用户表的user_id为1、2,订单表的user_id为1887、838,没有共同值,就算代码正确也不会输出匹配的关联结果

正确实现代码

使用Python内置的csv模块解析行,避免处理带引号的字段时出现拆分错误:

import apache_beam as beam
import csv
from io import StringIO

def parse_user_row(row):
    # 跳过表头行
    if row.startswith('user_id'):
        return
    reader = csv.reader(StringIO(row))
    fields = next(reader)
    user_id = int(fields[0])
    # 返回符合CoGroupByKey要求的KV结构
    return (user_id, tuple(fields[1:]))

def parse_order_row(row):
    # 跳过表头行
    if row.startswith('order_no'):
        return
    reader = csv.reader(StringIO(row))
    fields = next(reader)
    user_id = int(fields[1])
    return (user_id, tuple(fields[0:1] + fields[2:]))

if __name__ == "__main__":
    with beam.Pipeline() as pipeline:
        orders = (
            pipeline 
            | "Read orders" >> beam.io.ReadFromText("orders_v.csv")
            | "Parse orders to KV" >> beam.Map(parse_order_row)
            | "Filter invalid order rows" >> beam.Filter(lambda x: x is not None)
        )

        users = (
            pipeline 
            | "Read users" >> beam.io.ReadFromText("users_v.csv")
            | "Parse users to KV" >> beam.Map(parse_user_row)
            | "Filter invalid user rows" >> beam.Filter(lambda x: x is not None)
        )

        ({"orders": orders, "users": users} 
         | "Join by user_id" >> beam.CoGroupByKey()
         | "Print result" >> beam.Map(print)
        )

效果验证

修改orders_v.csv中任意一条订单的user_id为用户表存在的1或2,即可看到合并输出结果,示例输出如下:

(1, {'orders': [('1000', 'Cassava', '2000-01-01')], 'users': [('Anthony Wolf', 'male', '73', 'New Rachelburgh-VA-49583', '2019/03/13')]})

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 20:45:10