【问题标题】:Dataflow job does not any produce output数据流作业不产生任何输出
【发布时间】: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


    【解决方案1】:

    请问您指的是 2.22,因为“apache-beam 1.22”似乎已经过时了?特别是当您使用 Python 3.7 时,您可能想尝试更新的 SDK 版本,例如 2.22.0。

    我想要的是,如果它在窗口中接收到事件,无论如何它都应该在窗口之后产生输出。源是 Cloud PubSub,每分钟大约有 100 个事件。

    如果您只需要每个窗口触发一个窗格并每 5 分钟固定窗口一次,您可以简单地使用

    beam.WindowInto(beam.window.FixedWindows(5*60))
    

    如果你想自定义触发器,可以看看这个文档streaming-102。

    这是一个带有窗口输出可视化的流式示例。

    from apache_beam.runners.interactive import interactive_beam as ib
    
    ib.options.capture_duration = timedelta(seconds=30)
    ib.evict_captured_data()
    
    pstreaming = beam.Pipeline(InteractiveRunner(), options=options)
    
    words = (pstreaming
            | 'Read' >> beam.io.ReadFromPubSub(topic=topic)
            | 'Window' >> beam.WindowInto(beam.window.FixedWindows(5)))
    
    
    ib.show(words, visualize_data=True, include_window_info=True)
    

    如果您在 jupyterlab 等笔记本环境中运行这些代码,您可以使用this 等输出来调试流式传输管道。注意窗口是可视化的,在 30 秒内,我们得到 6 个窗口,因为固定窗口设置为 5 秒。您可以按窗口对数据进行分箱,以查看哪些数据来自哪个窗口。

    您可以按照instructions 设置自己的笔记本运行时; 或者您可以使用Google Dataflow Notebooks 提供的托管解决方案。

    【讨论】:

    • 谢谢,这是非常有价值的信息!没错,SDK的版本是2.22
    猜你喜欢
    • 2017-02-03
    • 2012-01-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-04-10
    相关资源
    最近更新 更多