【问题标题】:Run a function with no input inside beam.pipeline using Dataflow runner使用 Dataflow 运行器在 beam.pipeline 内运行没有输入的函数
【发布时间】:2021-11-19 15:14:41
【问题描述】:

我对 Dataflow 很陌生。我有一个批处理作业,可以从网站上抓取数据并将数据加载到 Bigquery

import apache_beam as beam
   

def get_data_from_id(id):
   # scrape data using input id
   return id

def process_data(data):
   # process data to dataframe
   return dataframe

def load_table(dataframe):
   # load dataframe to bq table
   # not return anything

# create beam pipeline 
p.beam.Pipeline()
id_no = p | "Input id" >> beam.create(['G123', 'G244', 'G444'])
data = id_no | "Scrape data" >> beam.Map(get_data_from_id)
data = data | "prcess data" >> beam.FlatMap(process_data)
data | "load to bq" >> beam.Map(load_table)

result = p.run()

这可以正常工作,并且能够使用数据流运行器将数据加载到 bq 表。 但我的问题是,我想在“加载到 bq”完成后添加另一个进程(作为普通的 python 函数),而不需要“加载到 bq”的输出结果

所以,我尝试了

def addtional_function():
   # run something 
   # no input argument 
...
...
result = p.run()
result.wait_until_finish()
additional_function()

使用 DirectRunner 运行良好,但我需要使用 Dataflow 运行程序运行它,我发现它仅适用于 beam.pipeline 内的函数。该函数未显示在 Dataflow 控制台中,也未运行。

附加函数不需要任何输入,但如果我将| "do additional fn" >> beam.Map(additional_function) 添加到光束流的下一行,则会出错,即有 1 个输入但附加函数()得到 0 个参数

那么,我如何将此函数包含到 beam.pipeline 中,以便进程能够使用 Dataflow runner 运行。

编辑:

如果我这样做

data = data | "load to bq" >> beam.Map(load_table)
data | "addtional fn" >> beam.Map(additional_function)

它会抛出这样的错误

TypeError: additional_function() takes 0 positional arguments but 1 was given [while running 'addtional fn']

如果我不这样解析输出,

data = data | "load to bq" >> beam.Map(load_table)
"addtional fn" >> beam.Map(additional_function)

additional_function() 也不会运行

【问题讨论】:

  • 你的新函数的输入是什么?与“加载到 bq”相同的输入?你不能使用分支吗?
  • @SergioLema 新函数不需要前一个函数的任何输入。只需运行它自己的函数
  • 如果函数不在 apache Beam 工作流中,则不会上传到 Dataflow。如果您希望此函数包含在 Datalfow 中,则必须在 wait_until_finish 之前将其包含在您的管道中
  • 是的,这就是问题所在。我应该如何做到这一点而不必解析 bean 管道中上一个进程的输出
  • 你的函数additional_function,应该作为一个map,但是你不会处理输入数据,忽略它。

标签: python google-cloud-dataflow apache-beam


【解决方案1】:

您的附加函数必须包含在 Beam 管道中才能导出到 Dataflow。

如果你希望它在最后执行,在“加载到 bq”之后,试试这个:

def additional_function(param):
   # all the logic you want
   return None # as it's a map function, you need an input parameter and an output value


data = data | "load to bq" >> beam.Map(load_table)
data | "additional_function" >> beam.Map(additional_function)

【讨论】:

  • 所以,即使我们不使用它,我们仍然需要输入吗?函数additional_function() 应该有一个输入参数?
  • 我的意思是,这只是一种解决方法,以防我们无法避免在附加函数()上输入参数。
  • 函数的签名必须有一个输入参数和一个返回值才能被接受为一个 Map 函数。但是在这种情况下你都不使用,所以忽略它们。
  • 谢谢。但只是好奇,这是处理这种情况的正确方法吗?
  • 是的,在函数式编程中是的
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-10-16
  • 1970-01-01
  • 2014-11-05
  • 1970-01-01
  • 2017-03-21
  • 2011-02-25
  • 1970-01-01
相关资源
最近更新 更多