【发布时间】:2021-11-20 10:28:03
【问题描述】:
我有两个 csv 文件,需要使用 beam(Python SDK)将它们合并到一个公共列上。文件如下所示:
users_v.csv
user_id,name,gender,age,address,date_joined
1,Anthony Wolf,male,73,New Rachelburgh-VA-49583,2019/03/13
2,James Armstrong,male,56,North Jillianfort-UT-86454,2020/11/06
orders_v.csv
order_no,user_id,product_list,date_purchased
1000,1887,Cassava,2000-01-01
1001,838,"Calabash, Water Spinach",2000-01-01
我尝试了以下似乎可行的方法(没有错误),但我无法使用 beam.Map(print) 查看生成的 PCollection:
import apache_beam as beam
with beam.Pipeline() as pipeline:
orders = p | "Read orders" >> beam.io.ReadFromText("orders_v.csv")
users = p | "Read users" >> beam.io.ReadFromText("users_v.csv")
{"orders": orders, "users": users} | beam.CoGroupByKey() | beam.Map(print)
如何打印生成的 PCollection?
【问题讨论】:
-
这是完整的代码吗?如果是这样,你不加入他们,你需要做一个 KV 然后 cogbk
-
是的,我只是从每个 csv 中获取了最上面的行。你能用一个例子详细说明吗?从梁文档看来,将使用
{"orders": orders, "users": users}创建一个 KV
标签: python apache-beam