【发布时间】:2019-07-30 00:04:00
【问题描述】:
我正在尝试遵循使用 Python SDK for Apache Beam on DataFlow 的流式管道的缓慢更改查找缓存 (https://cloud.google.com/blog/products/gcp/guide-to-common-cloud-dataflow-use-case-patterns-part-1) 的设计模式。
我们的查找缓存参考表位于 BigQuery 中,我们能够读取它并将其作为辅助输入传递给 ParDo 操作,但无论我们如何设置触发器/窗口,它都不会刷新。
class FilterAlertDoFn(beam.DoFn):
def process(self, element, alertlist):
print len(alertlist)
print alertlist
… # function logic
alert_input = (p | beam.io.Read(beam.io.BigQuerySource(query=ALERT_QUERY))
| ‘alert_side_input’ >> beam.WindowInto(
beam.window.GlobalWindows(),
trigger=trigger.RepeatedlyTrigger(trigger.AfterWatermark(
late=trigger.AfterCount(1)
)),
accumulation_mode=trigger.AccumulationMode.ACCUMULATING
)
| beam.Map(lambda elem: elem[‘SOMEKEY’])
)
...
main_input | ‘alerts’ >> beam.ParDo(FilterAlertDoFn(), beam.pvalue.AsList(alert_input))
根据此处的 I/O 页面 (https://beam.apache.org/documentation/io/built-in/),它说 Python SDK 仅支持 BigQuery Sink 的流式传输,这是否意味着 BQ 读取是有界源,因此无法在此方法中刷新?
尝试在源上设置非全局窗口会导致侧输入中的 PCollection 为空。
更新: 在尝试实施 Pablo 的回答建议的策略时,使用侧面输入的 ParDo 操作不会运行。
有一个输入源连接到两个输出,其中一个使用侧输入。 Non-SideInput 仍将到达其目的地,并且 SideInput 管道不会进入 FilterAlertDoFn()。
通过将侧输入替换为虚拟值,管道将进入函数。是不是在等待一个不存在的合适窗口?
使用与上面相同的 FilterAlertDoFn(),我的 side_input 和 call 现在看起来像这样:
def refresh_side_input(_):
query = 'select col from table'
client = bigquery.Client(project='gcp-project')
query_job = client.query(query)
return query_job.result()
trigger_input = ( p | 'alert_ref_trigger' >> beam.io.ReadFromPubSub(
subscription=known_args.trigger_subscription))
bigquery_side_input = beam.pvalue.AsSingleton((trigger_input
| beam.WindowInto(beam.window.GlobalWindows(),
trigger=trigger.Repeatedly(trigger.AfterCount(1)),
accumulation_mode=trigger.AccumulationMode.DISCARDING)
| beam.Map(refresh_side_input)
))
...
# Passing this as side input doesn't work
main_input | 'alerts' >> beam.ParDo(FilterAlertDoFn(), bigquery_side_input)
# Passing dummy variable as side input does work
main_input | 'alerts' >> beam.ParDo(FilterAlertDoFn(), [1])
我尝试了几个不同版本的 refresh_side_input(),它们在检查函数内部的返回时报告了预期的结果。
更新 2:
我对 Pablo 的代码做了一些小的修改,得到了相同的行为 - DoFn 永远不会执行。
在下面的示例中,每当我发布到 some_other_topic 时,我都会看到“in_load_conversion_data”,但在发布到 some_topic
时永远不会看到“in_DoFn”import apache_beam as beam
import apache_beam.transforms.window as window
from apache_beam.transforms import trigger
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import SetupOptions
from apache_beam.options.pipeline_options import StandardOptions
def load_my_conversion_data():
return {'EURUSD': 1.1, 'USDMXN': 4.4}
def load_conversion_data(_):
# I will suppose that these are currency conversions. E.g.
# {'EURUSD': 1.1, 'USDMXN' 20,}
print 'in_load_conversion_data'
return load_my_conversion_data()
class ConvertTo(beam.DoFn):
def __init__(self, target_currency):
self.target_currency = target_currency
def process(self, elm, rates):
print 'in_DoFn'
elm = elm.attributes
if elm['currency'] == self.target_currency:
yield elm
elif ' % s % s' % (elm['currency'], self.target_currency) in rates:
rate = rates[' % s % s' % (elm['currency'], self.target_currency)]
result = {}.update(elm).update({'currency': self.target_currency,
'value': elm['value']*rate})
yield result
else:
return # We drop that value
pipeline_options = PipelineOptions()
pipeline_options.view_as(StandardOptions).streaming = True
p = beam.Pipeline(options=pipeline_options)
some_topic = 'projects/some_project/topics/some_topic'
some_other_topic = 'projects/some_project/topics/some_other_topic'
with beam.Pipeline(options=pipeline_options) as p:
table_pcv = beam.pvalue.AsSingleton((
p
| 'some_other_topic' >> beam.io.ReadFromPubSub(topic=some_other_topic, with_attributes=True)
| 'some_other_window' >> beam.WindowInto(window.GlobalWindows(),
trigger=trigger.Repeatedly(trigger.AfterCount(1)),
accumulation_mode=trigger.AccumulationMode.DISCARDING)
| beam.Map(load_conversion_data)))
_ = (p | 'some_topic' >> beam.io.ReadFromPubSub(topic=some_topic)
| 'some_window' >> beam.WindowInto(window.FixedWindows(1))
| beam.ParDo(ConvertTo('USD'), rates=table_pcv))
【问题讨论】:
-
这个问题很有趣 - 我会尽量在明天之前回复您。
-
嘿 Pablo - 你有机会看看这个。我也对它感兴趣;-)
-
啊是的 - 抱歉耽搁了。一秒……
-
嗨@hulahoof,你能找到解决方案吗?我也面临与 python 相同的问题
标签: python streaming apache-beam dataflow