【问题标题】:Count distinct values in a stream pipeline计算流管道中的不同值
【发布时间】:2017-02-14 01:49:36
【问题描述】:

我有一个看起来像这样的管道

pipeline.apply(PubsubIO.read.subscription("some subscription"))
            .apply(Window.into(SlidingWindow.of(10 mins).every(20 seconds)
                            .triggering(AfterProcessingTime.pastFirstElementInPane()
                    .plusDelayOf(20 seconds))
                    .withAllowedLateness(Duration.ZERO)
                    .accumulatingFiredPanes()))
            .apply(RemoveDuplicates.create())
            .apply(Window.discardingFiredPanes()) // this is suggested in the warnings under https://cloud.google.com/dataflow/model/triggers#window-accumulation-modes
            .apply(Count.<String>globally().withoutDefaults())

此管道显着高估了不同的值(20 倍正常值)。最初,我怀疑默认触发器可能导致了这个问题。我已经调整以使用不允许迟到/丢弃触发的窗格/使用处理时间的触发器,所有这些都有类似的过度计数问题。

我也尝试过ApproximateUnique.globally:它在管道构建过程中失败,因为异常看起来像 Default values are not supported in Combine.globally() if the output PCollection is not windowed by GlobalWindows. 似乎没有办法将withoutDefaults 添加到它(就像我们对Count.globally 所做的那样)。

有没有推荐的方法在数据流/束流管道中以合理的精度执行COUNT(DISTINCT)

附:我正在使用 Java 数据流 SDK 1.9.0。

【问题讨论】:

    标签: google-cloud-dataflow dataflow


    【解决方案1】:

    您的代码看起来不错;它不应该多计。请注意,您将每个元素放入 30 个窗口中,因此如果您有一个不感知窗口的接收器(相当于折叠所有滑动窗口),您将期望元素的数量恰好是 30 倍。如果您可以显示更多的管道或您如何观察计数,那可能会有所帮助。

    除此之外,我对管道有一些建议:

    • 我建议将RemoveDuplicates 的触发器更改为AfterPane.elementCountAtLeast(1);这将以更低的延迟为您提供相同的结果,因为稍后到达的元素将没有影响。此触发器和您当前的触发器永远不会重复触发。因此,设置accumulatingFiredPanes() 还是discardingFiredPanes() 实际上并不重要。这很好,因为任何一个都不会与您的管道的其余部分一起工作。
    • 我会在Count 之前安装一个新触发器。原因有点技术性,但我会尝试描述它:
      • 在您当前的管道中,安装在那里的触发器(RemoveDuplicates 的触发器的“继续触发器”)记录第一个元素的到达时间,并等待它接收到在该时间或之前生成的所有元素处理时间,由上游工作人员测量。存在一些不确定性,因为它会影响本地处理时间和其他工作人员的处理时间。
      • 如果您采纳我的建议并将触发器切换为RemoveDuplicates,那么继续触发器将是AfterPane.elementCountAtLeast(1),因此它总是会尽快发出计数然后丢弃更多数据,这是非常错误的。李>

    【讨论】:

    • 感谢您的回答。我使用石墨来观察计数,我仍在调查那里可能出现的问题。对于您的第一点:AfterPane.elementCountAtLeast(1) 会倒数吗?鉴于此触发器将在至少有一个元素时发出窗口,所以RemoveDuplicate 大部分时间只会有一个元素?
    • 如果你使用了Repeatedly.forever(AfterPane.elementCountAtLeast(1)),那么你会收到很多只有一个元素的输出。但是只有AfterPane.elementCountAtLeast(1)),它会在第一次输出后关闭窗口,丢弃该键的其余输入。
    猜你喜欢
    • 2020-01-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多