【发布时间】:2020-09-26 13:02:06
【问题描述】:
for table_name, key_pair in relation_repl_key.items():
try:
with beam.Pipeline(options=PipelineOptions()) as p:
PCollection = p | "Reading from source database" >> relational_db.ReadFromDB(
source_config=source_config,
table_name=table_name,
query="SELECT {} FROM {}".format(
key_pair["col"],
table_name
)
)
side_input = bq(p, sideinput_bq_config, table_name, key_pair["repl_key"])
except RuntimeError:
pass
else:
PCollection | "Selecting updated rows" >> beam.ParDo(
KeyCheck(), beam.pvalue.AsSingleton(side_input)
)
finally:
load(PCollection, table_name, key_pair["primary_key"], key_pair["jsonb_col"])
我无法访问with 块之外的PCollection。在 finally: 内运行 PCollection | beam.Map(print) 不会返回任何内容。
【问题讨论】:
-
尝试单独定义 p
p = beam.Pipeline(options=PipelineOptions())
标签: google-cloud-dataflow apache-beam-io apache-beam