【问题标题】:Renaming Column names and creating new column names using apache beam使用 Apache Beam 重命名列名并创建新的列名
【发布时间】:2022-09-30 07:37:43
【问题描述】:

我有一个 CSV 文件,其中有 2 列名为.

我正在使用带有 direct_runner 的数据流。

我的用例首先将列名更改为姓名然后使用 PTransform 连接姓名和姓氏并生成一个新列员工姓名

代码 :

import apache_beam as beam

p2= beam.Pipeline()

def splitrow(element):
  return element.split(\',\')

demodata0=(
    
    p2
      |beam.io.ReadFromText(\'gs://demo/MOCK_DATA.csv\')
      |beam.Map(splitrow)
      |beam.Map(lambda element : ( element[0]+\" \"+element[1]))
      |beam.io.WriteToText(\'gs://demo/temp/output2\')

)

p2.run()

输入表:

first_name      last_name
John             Miller
Smith            scott

输出表:

name   surname   employee_name
john    Miller    John Miller
Smith   Scott     smith Scott

谢谢

  • 你的问题是什么?
  • 嘿@dnnshssm 我的问题是如何创建一个新列,甚至更改 apache 梁中的列名

标签: python-3.x google-cloud-platform google-cloud-dataflow apache-beam


【解决方案1】:

我以前从未在 beam 中使用过 CSV 文件,但我建议使用自定义 DoFn(请参阅 here)。它看起来像这样:

class EnrichCsvData(beam.DoFn):
  def process(self, element):
    output_pcoll = {}
    # i don't know if the inputs are strings, you might need to adjust the code if not
    output_pcoll["name"] = element[0]
    output_pcoll["surname"] = element[1]
    output_pcoll["employee_name"] = element[0] + element[1]
    
    return output_pcoll

然后在您的管道中调用它:

p2
  |beam.io.ReadFromText('gs://demo/MOCK_DATA.csv')
  |beam.Map(splitrow)
  |beam.ParDo(EnrichCsvData())
  |...

【讨论】:

  • 您好,非常感谢您的帮助。正如您所提到的,我使用自定义 DoFn 获得了所需的输出。
【解决方案2】:

当您有复杂的逻辑并且需要做一些繁重的工作时,创建自己的 DoFn 非常棒。如果您只需要选择一些列名称并且具有相对简单的定义,就像这里的情况一样,您可以使用 beam.Select() 来创建schemas

import apache_beam as beam

p2= beam.Pipeline()

def splitrow(element):
  return element.split(',')

demodata0=(
    
    p2
      |beam.io.ReadFromText('gs://demo/MOCK_DATA.csv')
      |beam.Map(splitrow)
      |beam.Select(name=lambda element: element[0],
                   surname=lambda element: element[1],
                   full_name=lambda element: element[0]+" "+element[1])
      |beam.io.WriteToText('gs://demo/temp/output2')

)

p2.run()

【讨论】:

    猜你喜欢
    • 2020-07-22
    • 2019-08-07
    • 1970-01-01
    • 1970-01-01
    • 2023-04-02
    • 1970-01-01
    • 1970-01-01
    • 2023-03-25
    • 1970-01-01
    相关资源
    最近更新 更多