【问题标题】:Filtering duplicates out of an infinite DataStream with windows使用 windows 从无限 DataStream 中过滤重复项
【发布时间】:2021-08-06 08:48:45
【问题描述】:

我想从无限 DataStream 中过滤掉 Flink 中的重复项。我知道重复项仅在很小的时间窗口(最多 10 秒)内出现。我发现了一种非常简单的有前途的方法here。但它不起作用。它使用键控 DataStream 并仅返回每个窗口的第一条消息。 这是我的窗口代码:

DataStream<Row> outputStream = inputStream
                .keyBy(new MyKeySelector())
                .window(SlidingProcessingTimeWindows.of(Time.seconds(10), Time.minutes(5)))
                .process(new DuplicateFilter());

MyKeySelector()只是一个选择Row消息的前两个属性作为key的类。此键用作主键,并导致只有具有相同键的消息才会分配给同一窗口(经典键控流行为)。

这就是Duplicate Filter 类,它与上述问题的建议答案非常相似。我只使用了较新的process() 函数而不是apply()

public class DuplicateFilter extends ProcessWindowFunction<Row, Row, Tuple2<String, String>, TimeWindow> {
private static final Logger LOG = LoggerFactory.getLogger(DuplicateFilter.class);

@Override
public void process(Tuple2<String, String> key, Context context, Iterable<Row> iterable, Collector<Row> collector) throws Exception {
    // this is just for debugging and can be ignored
    int count = 0;
    for (Row record :
            iterable) {
        LOG.info("Row number {}: {}", count, record);
        count++;
    }
    LOG.info("first Row: {}", iterable.iterator().next());

    collector.collect(iterable.iterator().next()); //output only the first message in this window
}
}

我的消息以最大间隔到达。一秒钟,所以一个 30 秒的窗口应该能很好地处理这个问题。但是距离小于 1 秒的消息被分配到不同的窗口。我从日志中可以看出,它很少能正常工作。

有人对此任务有想法或其他方法吗?如果您需要更多信息,请告诉我。

【问题讨论】:

    标签: java apache-flink flink-streaming


    【解决方案1】:

    Flink 的时间窗口与时钟对齐,而不是与事件对齐,因此可以将时间上接近的两个事件分配到不同的窗口。 Windows 通常不太适合重复数据删除,但如果使用会话窗口,您可能会得到很好的结果。

    就我个人而言,我会使用键控平面图(或进程函数),并在不再需要键时使用状态 TTL(或计时器)来清除键的状态。

    您还可以使用 Flink SQL 进行重复数据删除:https://ci.apache.org/projects/flink/flink-docs-stable/docs/dev/table/sql/queries/deduplication/(但您需要 set an idle state retention interval)。


    【讨论】:

    • 谢谢。你能给出一个带有 TTL 的小代码示例吗? SessionWindows 没有解决问题。
    • 我已经通过将 DataStream 转换为表格来解决它。使用 SQL 进行重复数据删除非常简单。我不需要设置空闲状态保留间隔。除了给定的链接,我发现这个链接有助于将 DataStream 转换为表:ci.apache.org/projects/flink/flink-docs-release-1.13/docs/dev/…
    • 没有什么会强制您设置空闲状态保留间隔,但如果您不这样做,您最终可能会耗尽内存或磁盘空间——因为默认情况下重复数据删除的范围是无限的。
    猜你喜欢
    • 2018-12-09
    • 1970-01-01
    • 1970-01-01
    • 2021-09-27
    • 2012-10-18
    • 2014-07-10
    • 2016-02-12
    • 2016-06-06
    • 2017-03-18
    相关资源
    最近更新 更多