【问题标题】:How to merge two files and then view the PCollection (Apache Beam)如何合并两个文件然后查看 PCollection (Apache Beam)
【发布时间】: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


【解决方案1】:

代码有几个错误:

1 - 您在with 中使用pipeline,但随后您使用p 作为管道变量

2 - CoGroupByKey 之前的字典确定了组合变量的名称,但它仍然需要一个键值来加入

3 - 我猜你想跳过标题。

代码应如下所示。 split_by_kv 函数远非完美,您需要对其进行改进,以便更好地检索密钥(因为您的某些字段可能包含 ,)。

def split_by_kv(element, index, delimiter=", "):
    # Need a better approach here
    splitted = element.split(delimiter)
    return splitted[index], element


with beam.Pipeline() as p:
    orders = (p | "Read orders" >> ReadFromText("files/orders_v.csv", skip_header_lines=1)
                | "to KV order" >> Map(split_by_kv, index=1, delimiter=",")
             )
    
    users = (p | "Read users" >> ReadFromText("files/users_v.csv", skip_header_lines=1)
               | "to KV users" >> Map(split_by_kv, index=0, delimiter=",")
            )
    
    ({"orders": orders, "users": users} | CoGroupByKey() 
                                        | Map(print)
    )

输出是 (key, {"orders": values for key, "users": values for users})

('1887', {'orders': ['1000,1887,Cassava,2000-01-01'], 'users': []})
('838', {'orders': ['1001,838,"Calabash, Water Spinach",2000-01-01'], 'users': []})
('1', {'orders': [], 'users': ['1,Anthony Wolf,male,73,New Rachelburgh-VA-49583,2019/03/13']})
('2', {'orders': [], 'users': ['2,James Armstrong,male,56,North Jillianfort-UT-86454,2020/11/06']})

另外,您可能想看看新的DataFrames API

【讨论】:

    猜你喜欢
    • 2020-01-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-09
    • 2022-12-31
    • 2023-02-03
    • 2023-04-10
    • 2023-04-10
    相关资源
    最近更新 更多