【发布时间】: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