【发布时间】:2020-06-11 09:25:40
【问题描述】:
我有一个问题,数据流作业实际上运行良好,但在手动排空作业之前它不会产生任何输出。
使用以下代码,我假设它会产生窗口输出,有效地在每个窗口之后触发。
lines = (
p
| "read" >> source
| "decode" >> beam.Map(decode_message)
| "Parse" >> beam.Map(parse_json)
| beam.WindowInto(
beam.window.FixedWindows(5*60),
trigger=beam.trigger.Repeatedly(beam.trigger.AfterProcessingTime(5*60)),
accumulation_mode=beam.trigger.AccumulationMode.DISCARDING)
| "write" >> sink
)
我想要的是,如果它在窗口中接收到事件,无论如何它都应该在窗口之后产生输出。源是 Cloud PubSub,每分钟大约有 100 个事件。
这是我用来启动工作的参数:
python main.py \
--region $(REGION) \
--service_account_email $(SERVICE_ACCOUNT_EMAIL_TEST) \
--staging_location gs://$(BUCKET_NAME_TEST)/beam_stage/ \
--project $(TEST_PROJECT_ID) \
--inputTopic $(TOPIC_TEST) \
--outputLocation gs://$(BUCKET_NAME_TEST)/beam_output/ \
--streaming \
--runner DataflowRunner \
--temp_location gs://$(BUCKET_NAME_TEST)/beam_temp/ \
--experiments=allow_non_updatable_job \
--disk_size_gb=200 \
--machine_type=n1-standard-2 \
--job_name $(DATAFLOW_JOB_NAME)
关于如何解决这个问题的任何想法?我正在使用 apache-beam 2.22 SDK,python 3.7
【问题讨论】:
标签: python-3.x google-cloud-platform google-cloud-dataflow dataflow