【问题标题】:Beam/Dataflow stateful processing, ParDo never runsBeam/Dataflow 状态处理,ParDo 从不运行
【发布时间】:2020-03-31 15:24:02
【问题描述】:

我正在尝试在 Dataflow 上使用 Beam 的有状态处理,但每次尝试输出数据时,我都会在日志中收到这些错误。结果是有状态的ParDo+DoFn什么都没有输出:

16:45:56.948 CEST Proposing dynamic split of work unit myproject;2020-03-31_07_34_07-7523868393961495218;8536385410242733529 at {"fractionConsumed":0.5}
16:45:56.948 CEST Rejecting split request because custom reader returned null residual source.

编辑这似乎是巧合。 似乎有状态的ParDo 在窗口触发之前不会输出任何元素。这是正确的吗?

此示例使用 Scio 的 .batchByKey 复制错误(它在后台使用有状态处理):

    val create = Create.of(()).withCoder(CoderMaterializer.beam(sc, Coder[Unit]))
    sc.customInput("Unit input", create)
      .map(_ => println("STARTING"))
      .applyTransform(ParDo.of(new Increasing)) // Outputs infinite stream of increasing numbers, one per second, prints each number to stdout
      .keyBy(1 -> _)
      .batchByKey(5)
      .map {
        case (key, vs) => vs.foreach(v => println(s"GOT batch with $v"))
      }
    sc.run()

最终的.map,它只是一个带有单个输出的ParDo+DoFn,永远不会运行。

在输出中,我看到五行不断增加的数字(来自new Increasing),然后是上面的两条消息。如此反复。

有人知道错误可能是什么吗?这似乎是来源apache/beam/../WorkerCustomSources.java#L698

【问题讨论】:

    标签: google-cloud-dataflow apache-beam spotify-scio


    【解决方案1】:

    该错误仅表示您在管道中使用的源无法拆分。因此,它与您的管道进度没有直接关系,但在这种情况下可能与此相关。我看到您使用的唯一来源是Create,您的问题可能是由于您的Create 使用空元组初始化。可以尝试使用具有一个或多个元素的 Create 为您的管道播种。

    【讨论】:

    • 同意,这似乎是巧合。你知道有状态处理与windows有什么关系吗?有状态的 ParDo 只会在窗口触发时输出吗?
    • 我认为触发取决于您设置的窗口/触发机制。与您的 ParDo 是否有状态无关。我认为这里的问题可能是您的 Create 转换为空。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-25
    • 2021-03-11
    • 2021-08-13
    • 1970-01-01
    相关资源
    最近更新 更多