【发布时间】: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