【问题标题】:Cloud dataflow Watermark stuck and increasing system lag云数据流水印卡住并增加系统延迟
【发布时间】:2019-06-18 06:39:48
【问题描述】:

我正在从数据流管道中的 PubSub 主题读取记录。 PubSub 记录分为固定窗口,然后在每个窗口上分组。每个窗口都按序列号排序,因为我们需要使用 beam.SortValues 按顺序处理这些记录。然后我将记录写入 Cloud BigTable

管道的问题是数据新鲜度和系统滞后。数据新鲜度似乎卡在某个点,水印停止前进。

我正在使用以下窗口策略在 GroupByKey 步骤之后发出记录:

PCollection<KV<BigInteger, JSONObject>> window = pubsubRecords.apply("Raw to String", ParDo.of(new LogsFn()))
                .apply("Window", Window
                        .<KV<BigInteger, JSONObject>>into(FixedWindows.of(Duration.standardSeconds(10)))
                        .triggering(Repeatedly.forever(AfterFirst.of(
      AfterPane.elementCountAtLeast(500),
      AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1)))))
                        .withAllowedLateness(Duration.ZERO).discardingFiredPanes()
                    );

我认为问题可能出在窗口策略上。基本上我想做以下事情:从 PubSub 读取记录到 1 分钟的 FixedWindows,对窗口进行排序并写入 BigTable。如果我使用默认触发器,则 GroupByKey 步骤不会发出任何结果。有人可以帮我解决这个问题吗?

【问题讨论】:

    标签: java streaming google-cloud-dataflow apache-beam


    【解决方案1】:

    阅读您的代码,现在看起来您的早期触发器和窗口大小是向后的。您的窗口策略实际上是:

    1. 10 秒事件时间固定窗口
    2. 1 分钟处理时间或窗格中 500 个元素的复合提前触发。
    3. 迟到的事件被丢弃。

    如果您只想要 1 分钟的事件时间窗口,这就是您所需要的:

    PCollection<KV<BigInteger, JSONObject>> window = pubsubRecords.apply("Raw to String", ParDo.of(new LogsFn()))
                .apply("Window", Window
                .<KV<BigInteger, JSONObject>>into(FixedWindows.of(Duration.standardMinutes(1)))
                    .withAllowedLateness(Duration.ZERO)
                    .discardingFiredPanes()
                    .withOnTimeBehavior(OnTimeBehavior.FIRE_ALWAYS));
    

    Fire 始终是默认的 OnTimeBehavior,但为了便于阅读,我们可以使其明确。如果您需要复合触发器,可以将其重新添加 - 我怀疑您想触发一个 10 秒或 500 个元素。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-06
      • 1970-01-01
      • 1970-01-01
      • 2017-10-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多