【问题标题】:Dataflow Distinct transform example数据流不同的转换示例
【发布时间】: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


【解决方案1】:

因为您想在 1 分钟窗口内删除重复项;

您可以使用带有默认触发器的固定窗口,而不是使用带有处理时间触发器的全局窗口。

Window.<String>into(FixedWindows.of(Duration.standardSeconds(60))));

随后的独特转换将根据事件时间删除 1 分钟窗口内的所有重复键。

【讨论】:

  • 谢谢。愚蠢的问题:默认触发事件时间是驱动的吗? guide 似乎表明了这一点。我需要一个可预测的时间来“处理跨窗口的重复项”,除非我引入其他内容,例如允许唯一列值的数据库。因此,我可以根据处理时间来执行此操作吗?
  • 固定窗口将是可预测的:beam.apache.org/documentation/programming-guide/… 对于每个 1 分钟的固定窗口,元素将被重复数据删除。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2016-11-30
  • 2020-05-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多