【发布时间】:2019-12-15 10:59:40
【问题描述】:
在我的 Dataflow 管道中,我尝试使用 Distinct transform 来减少重复项。我想尝试最初将其应用于固定的 1 分钟窗口,并使用另一种方法来处理窗口中的重复项。如果 1 分钟的窗口是真实的/处理时间,则后一点可能效果最好。
我希望有 1000 个元素,每个文本字符串只有几 KiB。
我这样设置 Window 和 Distinct 变换:
PCollection<String>.apply("Deduplication global window", Window
.<String>into(new GlobalWindows())
.triggering(Repeatedly
.forever(AfterProcessingTime
.pastFirstElementInPane()
.plusDelayOf(Duration.standardMinutes(1))
)
)
.withAllowedLateness(Duration.ZERO).discardingFiredPanes()
)
.apply("Deduplicate URLs in window", Distinct.<String>create());
但是当我在 GCP 上运行它时,我看到 Distinct 转换发出的元素似乎比它接收的要多:
(因此,根据定义,除非它构成了某些东西,否则它们不可能是不同的!)
更有可能我猜我没有正确设置它。有没有人有一个如何做到这一点的例子(除了javadoc,我真的没有找到太多)?谢谢。
【问题讨论】:
-
您是否试图在管道的整个生命周期中找到不同的元素?或者您希望每个键在特定时间范围内具有不同的元素?
-
我希望在每个处理时间窗口(1 分钟)内有不同的元素 - 我不想累积比这更大的状态。
-
我们为此放弃了Dataflow,但从那以后我了解到额外的输出元素通常是由于一批元素中的错误,导致整个批次被重新处理和重新计算,所以实际上,就唯一性而言,输出集合的大小并不意味着什么。
标签: java apache-beam dataflow distinct-values