我將side_inputPCollection 作為側面輸入傳遞給ParDo轉換,但得到相同的 KeyErrorimport apache_beam as beamfrom apache_beam.options.pipeline_options import PipelineOptionsfrom beam_nuggets.io import relational_dbfrom processors.appendcol import AppendColfrom side_inputs.config import sideinput_bq_configfrom source.config import source_configwith beam.Pipeline(options=PipelineOptions()) as si: side_input = si | "Reading from BQ side input" >> relational_db.ReadFromDB( source_config=sideinput_bq_config, table_name='abc', query="SELECT * FROM abc" )with beam.Pipeline(options=PipelineOptions()) as p: PCollection = p | "Reading records from database" >> relational_db.ReadFromDB( source_config=source_config, table_name='xyzzy', query="SELECT * FROM xyzzy", ) | beam.ParDo( AppendCol(), beam.pvalue.AsIter(side_input) )下面是錯誤Traceback (most recent call last): File "athena/etl.py", line 40, in <module> extract() File "athena/etl.py", line 22, in extract PCollection = p | "Reading records from database" >> relational_db.ReadFromDB( File "/Users/souvikdey/.pyenv/versions/3.8.5/envs/athena-venv/lib/python3.8/site-packages/apache_beam/pipeline.py", line 555, in __exit__ self.result = self.run() File "/Users/souvikdey/.pyenv/versions/3.8.5/envs/athena-venv/lib/python3.8/site-packages/apache_beam/pipeline.py", line 534, in run return self.runner.run_pipeline(self, self._options) File "/Users/souvikdey/.pyenv/versions/3.8.5/envs/athena-venv/lib/python3.8/site-packages/apache_beam/runners/direct/direct_runner.py", line 119, in run_pipeline return runner.run_pipeline(pipeline, options) File "/Users/souvikdey/.pyenv/versions/3.8.5/envs/athena-venv/lib/python3.8/site-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 175, in run_pipeline self._latest_run_result = self.run_via_runner_api(我正在從 PostgreSQL 表中讀取數據,PCollection 的每個元素都是一個字典。
1 回答

猛跑小豬
TA貢獻1858條經驗 獲得超8個贊
我認為問題在于你有兩個獨立的管道試圖一起工作。您應該將所有轉換作為單個管道的一部分執行:
with beam.Pipeline(options=PipelineOptions()) as p:
side_input = p | "Reading from BQ side input" >> relational_db.ReadFromDB(
source_config=sideinput_bq_config,
table_name='abc',
query="SELECT * FROM abc")
my_pcoll = p | "Reading records from database" >> relational_db.ReadFromDB(
source_config=source_config,
table_name='xyzzy',
query="SELECT * FROM xyzzy",
) | beam.ParDo(
AppendCol(), beam.pvalue.AsIter(side_input))
添加回答
舉報
0/150
提交
取消