如何在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
相关产品推荐
相关产品推荐

