Advertisement

FlinkDataStream如何实现去重

阅读量:

摘要

假设有批订单数据实时接入Kafka系统, 而Flink系统则需要对这些实时收到的订单数据进行处理。特别注意的是,在实际应用中不允许对相同的订单数据进行重复处理。为此,在Flink系统的具体流程设计中必须确保不会出现这种情况。具体而言,在FlinkSQL框架内已经提供了相应的去重功能,请参考FlinkSQL去重一文中的详细说明。然而,在本节我们将重点讨论DataStream API下的去重机制及其实现方式。

分析

我们很容易想到:假设订单的唯一主键就是order_id, 要想达到去重的效果应该可以想到用State 存储已经处理过的订单,新的订单来临的时候判断是否存在于State中,如果不存在则处理,存在则视为重复订单,需要放弃当前订单。
上面的思想理论上是没有问题的,但是实际上却会产生不小的问题。 上面的额分析中,state会缓存所有已经处理过的订单id, 要知道kakfa的数据是源源不断的, 那么也就意味着我们需要缓存的state 会越来越大, 没错这就像一个不断膨胀的炸弹,总有一天会炸掉。因此我们需要在分析下。 也就是说Datastream 的缓存(状态)不能一直存在,否则总会内存溢出。 如果是你你会怎么做? 没错给缓存增加一个过期时间。 而这个时间就要结合业务来确定。
加入我有订单数据, 正常来说即便是重复订单,一定会在一个小时之内过

全部评论 (0)

还没有任何评论哟~