Flink 实时去重三部曲:HashSet → 布隆过滤器 → 布隆 + 窗口触发器渐进式优化实战
一、前言
实时数据流中重复数据是十分常见的问题:上游 MQ 重试、网络抖动、业务重发都会产生重复消息,如果不做去重处理,会导致统计指标虚高、入库数据冗余、业务对账出错。 针对 Flink 实时去重需求,本文循序渐进实现三套方案,由浅入深分析优缺点、内存开销、适用场景:
- 方案一:内存
HashSet全量存储去重(最简单入门版) - 方案二:布隆过滤器 BloomFilter 优化内存占用(工业常用优化版)
- 方案三:布隆过滤器 + 窗口触发器定时清理(解决内存无限膨胀生产可用版)
二、业务背景
当前智慧交通需求:按区域 ID 分组,统计 10 分钟滚动事件时间窗口内,各个区域的独立过车数量。同一辆车短时间多次抓拍会产生重复数据,必须对车牌做去重统计,避免车流量统计虚高。
三、基于 HashSet 内存去重(初级实现)
核心实现思路
- 按照
areaId分组,开启 10 分钟事件时间滚动窗口; - 遍历窗口内所有车辆数据,将车牌存入
HashSet自动剔除重复车牌; - 集合最终大小 = 当前区域窗口内去重后的真实车辆总数;
- 封装统计结果写入 MySQL。
核心代码
ds2.keyBy(CarInfo::getAreaId)
.window(TumblingEventTimeWindows.of(Time.minutes(10)))
.apply(new WindowFunction<CarInfo, AreaControl, String, TimeWindow>() {
@Override
public void apply(String key, TimeWindow window, Iterable<CarInfo> input, Collector<AreaControl> out) throws Exception {
// 定义HashSet存储车牌,自动去重
HashSet<String> carPlateSet = new HashSet<>();
// 遍历窗口全部车辆数据
for (CarInfo carInfo : input) {
carPlateSet.add(carInfo.getCar());
}
// 封装统计结果
AreaControl areaControl = new AreaControl();
areaControl.setAreaId(key);
areaControl.setWindowStart(DateFormatUtils.format(window.getStart(), "yyyy-MM-dd HH:mm:ss"));
areaControl.setWindowEnd(DateFormatUtils.format(window.getEnd(), "yyyy-MM-dd HH:mm:ss"));
// 集合长度就是去重后独立车辆数
areaControl.setCarCount(carPlateSet.size());
out.collect(areaControl);
}
});
方案优缺点分析
✅ 优点
- 实现逻辑简单易懂,无任何误判,精准去重;
- 不需要引入第三方依赖,JDK 原生集合即可实现;
- 调试方便,适合理解窗口内聚合去重基础原理。
❌ 缺点
- 高车流场景下,窗口内车牌数量巨大,
HashSet存储完整字符串,堆内存占用持续走高,存在 OOM 内存溢出风险; - 窗口迭代器会一次性加载当前窗口所有数据到内存,大数据量吞吐下 GC 压力大;
- 仅能实现单窗口内局部去重,无法支撑跨窗口全局去重场景;
- 窗口触发前所有数据常驻内存,高峰期资源开销明显,仅适合测试、小流量场景,不建议直接上生产。
存在痛点引出后续优化
为解决 HashSet 内存占用过高问题,我们引入空间利用率更高的布隆过滤器,作为第二版优化方案。
四、Redis 布隆过滤器窗口去重(内存优化版)
什么是布隆过滤器(Bloom Filter)
布隆过滤器是空间利用率极高的概率型二进制数据结构,底层本质是一个超长 bit 位数组,数组内元素只有 0、1 两种值。
核心判定特性
- 如果判定元素一定不存在:返回结果百分百准确;
- 如果判定元素可能存在:存在极小概率哈希碰撞造成误判(不存在的元素被误认为已存在);
- 不存储原始数据本身,只通过多个哈希函数映射标记 bit 位,内存占用远小于 HashSet、Redis Set。
基础插入 & 查询逻辑
- 插入元素:对同一个车牌执行多个不同哈希运算,算出多个 bit 下标,把对应位置全部置为
1; - 判断重复:取出多个哈希下标,校验所有 bit 位是否全为 1;只要有任意一位是 0,代表车牌从未出现;全部为 1 则判定为已存在。
布隆过滤器与 Redis 的关系
- Redis 原生提供
setbit、getbit指令,可以直接操作内部二进制位图(BitMap),完美实现布隆过滤器底层位数组; - 我们以区域 ID + 窗口起始时间作为 Redis Key,隔离不同区域、不同窗口的去重数据,避免上一个窗口数据干扰当前窗口统计;
- 对比方案:
- HashSet:存储完整车牌字符串,海量数据堆内存暴涨,容易 OOM;
- Redis Set:存储完整车牌字符串,海量场景内存开销大、成本高;
- Redis BitMap 实现布隆过滤器:仅用 bit 位标记存在性,存储空间压缩几十上百倍,适合卡口千万级车辆实时去重。
如何降低布隆过滤器误判率
布隆过滤器无法彻底消除误判,只能通过配置压低误判概率,本项目采用两种优化手段:
- 增加哈希函数数量 本案例使用
car.hashCode()+MD5哈希两套独立哈希算法映射 bit 位,哈希函数越多,碰撞概率越低; - 扩大 bit 数组总长度 预设 bit 总长度
length = 5000000,数组容量越大,哈希下标碰撞概率越低; - 业务适配:交通抓拍重复本就是短时间内高频重试,极低误判率完全可以容忍,满足业务统计指标要求。
Redis + 布隆过滤器 完整核心代码
哈希工具类 BloomFilterUtil
package com.bigdata.utils;
import java.math.BigInteger;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.Arrays;
public class BloomFilterUtil {
// 布隆过滤器bit数组总长度
private static int length = 5000000;
// MD5哈希算法
public static Integer md5Hash(String input) {
try {
MessageDigest digest = MessageDigest.getInstance("MD5");
byte[] hashBytes = digest.digest(input.getBytes(StandardCharsets.UTF_8));
BigInteger hashNumber = new BigInteger(1, hashBytes);
return hashNumber.intValue();
} catch (NoSuchAlgorithmException e) {
e.printStackTrace();
return null;
}
}
// 生成两个哈希偏移下标
public static int[] getOffsets(String car) {
int[] arr = new int[2];
// 哈希1:字符串原生hashCode
int a = Math.abs(car.hashCode()) % length;
// 哈希2:MD5自定义哈希
int b = Math.abs(md5Hash(car)) % length;
arr[0] = a;
arr[1] = b;
return arr;
}
}
两种hash算法
- 哈希 1:
字符串原生hashCode → 取绝对值 → 对数组长度取模 = 下标1 - 哈希 2:
字符串MD5摘要转整数 → 取绝对值 → 对数组长度取模 = 下标2
Flink 窗口去重核心业务代码(仅窗口逻辑)
SingleOutputStreamOperator<AreaControl> resultStream = keyedStream
.window(TumblingEventTimeWindows.of(Time.minutes(10)))
.apply(new WindowFunction<CarInfo, AreaControl, String, TimeWindow>() {
@Override
public void apply(String areaId, TimeWindow window, Iterable<CarInfo> input, Collector<AreaControl> out) throws Exception {
AreaControl areaControl = new AreaControl();
areaControl.setAreaId(areaId);
long start = window.getStart();
long end = window.getEnd();
String startStr = DateFormatUtils.format(start, "yyyy-MM-dd HH:mm:ss");
String endStr = DateFormatUtils.format(end, "yyyy-MM-dd HH:mm:ss");
areaControl.setWindowStart(startStr);
areaControl.setWindowEnd(endStr);
// 定义Redis Key:区域ID+窗口起始时间,隔离不同窗口数据
String redisKey = areaId + ":" + startStr;
Jedis jedis = new Jedis("localhost", 6379);
int carCount = 0;
for (CarInfo carInfo : input) {
String carPlate = carInfo.getCar();
// 获取两个哈希偏移位置
int[] offsets = BloomFilterUtil.getOffsets(carPlate);
Boolean bit1 = jedis.getbit(redisKey, offsets[0]);
Boolean bit2 = jedis.getbit(redisKey, offsets[1]);
// 任意一个bit位为0,代表车牌不存在,统计+1并置位
if (!bit1 || !bit2) {
carCount++;
jedis.setbit(redisKey, offsets[0], true);
jedis.setbit(redisKey, offsets[1], true);
}
}
areaControl.setCarCount(carCount);
out.collect(areaControl);
jedis.close();
}
});
判断存在逻辑
得到一个车牌号,然后根据对应的rediskey去redis中去查找,如果任意一个bit为flase则车牌不存在,计数器加1,并且在对应的redis中的位置设置为true。
窗口过期清理补充说明
使用 areaId+窗口起始时间 作为 Redis Key,天然隔离不同窗口数据;可配置 Redis 过期策略,窗口结束一段时间后自动删除对应 key,避免 Redis 数据长期堆积。
方案优缺点分析
✅ 优点
- 存储空间极度精简,相比 HashSet、Redis Set 大幅节约内存 / Redis 资源,支撑卡口海量车辆去重统计;
- 借助 Redis 外置存储,Flink 任务重启不会丢失窗口去重标记,稳定性优于内存 HashSet;
- 双哈希设计压低误判率,满足交通车流量统计业务容忍度;
- 适合大流量实时去重场景,解决 HashSet 大流量 OOM 致命问题。
❌ 缺点
- 存在极低概率哈希碰撞误判,无法做到 100% 精准去重,对零误差强一致性业务不适用;
- 频繁创建关闭 Jedis 连接,高并发下存在网络 IO 开销,可改造连接池优化;
- 布隆过滤器不支持删除单个元素,只能整体删除整个 Redis Key;
- 只能实现窗口内局部去重,无法天然实现跨窗口全局去重。
本方案遗留痛点(引出方案三)
当前 Redis 布隆过滤器绑定固定窗口 Key,虽然解决内存溢出问题,但需要手动管理 Redis 过期清理;如果业务窗口数量极多,Redis Key 会持续累积、运维繁琐。 下一步优化思路:布隆过滤器 + Flink 窗口触发器,在窗口触发结束时自动清理当前窗口布隆过滤器数据,实现生命周期全自动管理,形成生产级稳定去重方案。
五、布隆过滤器 + 自定义触发器 逐条实时去重(生产最终优化方案)
现有前两套方案遗留核心痛点回顾
- 方案一 HashSet:窗口等待所有数据齐了再遍历迭代器统计,超大流量下迭代器积攒海量对象,堆内存暴涨极易 OOM;
- 方案二 Redis 布隆过滤器:依然依赖窗口攒齐全部数据再统一遍历计算,迭代器数据积压问题没有解决;同时逐条落地 MySQL 频繁创建关闭连接,IO 压力大、写入性能极差。
针对性优化思路: 自定义窗口触发器,每来一条数据立刻触发窗口计算、用完立即清空窗口数据,不让迭代器堆积数据;统计结果落地 Redis 替代 MySQL,解决频繁入库性能瓶颈,结合 Redis 布隆过滤器实现精准去重统计。
Flink Window 触发器(Trigger)原理详解
触发器作用
Trigger 是窗口的调度控制器,用来决定什么时候触发窗口函数执行计算。 默认滚动 / 滑动窗口只会在窗口结束水位线到达后才触发一次计算;我们自定义触发器可以改写触发时机,实现来一条数据就执行一次窗口逻辑。
Trigger 四个核心重写方法
onElement():窗口每流入一条数据,就会执行一次onEventTime():事件时间定时器触发时执行onProcessingTime():处理时间定时器触发时执行clear():窗口销毁时,做资源清理工作
TriggerResult 四种返回枚举(控制窗口行为)
CONTINUE:不做任何操作,继续等待数据FIRE:执行窗口计算输出结果,保留窗口内原有数据PURGE:直接清空窗口所有数据,不执行计算、不输出FIRE_AND_PURGE:先执行窗口计算输出结果,计算完成立刻清空窗口所有数据(本方案核心使用)
自定义触发器完整代码
package com.bigdata;
import com.bigdata.pojo.CarInfo;
import org.apache.flink.streaming.api.windowing.triggers.Trigger;
import org.apache.flink.streaming.api.windowing.triggers.TriggerResult;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
public class MyTrigger extends Trigger<CarInfo, TimeWindow> {
// 每条数据进入窗口都会执行
@Override
public TriggerResult onElement(CarInfo element, long timestamp, TimeWindow window, TriggerContext ctx) throws Exception {
// 来一条计算一条,计算后清空窗口,防止数据积压
return TriggerResult.FIRE_AND_PURGE;
}
@Override
public TriggerResult onProcessingTime(long time, TimeWindow window, TriggerContext ctx) throws Exception {
return TriggerResult.CONTINUE;
}
@Override
public TriggerResult onEventTime(long time, TimeWindow window, TriggerContext ctx) throws Exception {
return TriggerResult.CONTINUE;
}
// 窗口销毁时回调清理资源
@Override
public void clear(TimeWindow window, TriggerContext ctx) throws Exception {
}
}
关键逻辑:onElement 返回 FIRE_AND_PURGE,实现单条数据触发计算 + 即时清空窗口,迭代器永远不会堆积多条数据,从根源解决大流量内存溢出问题。
核心代码(仅窗口 + 业务逻辑精简版)
SingleOutputStreamOperator<AreaControl> resultStream = keyedStream
.window(TumblingEventTimeWindows.of(Time.minutes(10)))
//触发器在此定义
.trigger(new MyTrigger())
.apply(new WindowFunction<CarInfo, AreaControl, String, TimeWindow>() {
// 维护区域+窗口维度车辆统计总数
HashMap<String,Integer> hashMap=new HashMap<>();
@Override
public void apply(String areaId, TimeWindow window, Iterable<CarInfo> input, Collector<AreaControl> out) throws Exception {
long start = window.getStart();
String startStr = DateFormatUtils.format(start, "yyyy-MM-dd HH:mm:ss");
String endStr = DateFormatUtils.format(window.getEnd(), "yyyy-MM-dd HH:mm:ss");
String mapKey= "areaId:"+areaId+",startTime:"+startStr;
String bloomKey = areaId+":"+startStr;
Jedis jedis = new Jedis("localhost", 6379);
for (CarInfo carInfo : input) {
String car= carInfo.getCar();
int[] offsets = BloomFilterUtil.getOffsets(car);
Boolean k1 = jedis.getbit(bloomKey, offsets[0]);
Boolean k2 = jedis.getbit(bloomKey, offsets[1]);
if(!hashMap.containsKey(mapKey)){
hashMap.put(mapKey,1);
jedis.setbit(bloomKey,offsets[0],true);
jedis.setbit(bloomKey,offsets[1],true);
}else{
Integer carNum = hashMap.get(mapKey);
if(!k1 || !k2){
hashMap.put(mapKey,carNum+1);
jedis.setbit(bloomKey,offsets[0],true);
jedis.setbit(bloomKey,offsets[1],true);
}
}
}
AreaControl areaControl = new AreaControl();
areaControl.setAreaId(areaId);
areaControl.setWindowStart(startStr);
areaControl.setWindowEnd(endStr);
areaControl.setCarCount(hashMap.get(mapKey));
out.collect(areaControl);
jedis.close();
}
});
代码逻辑简述
-
窗口与触发机制 基于10 分钟事件时间滚动窗口,绑定自定义触发器
MyTrigger,每流入一条车辆数据就立刻执行窗口计算,计算完成马上清空窗口,避免窗口迭代器积攒大量数据导致内存溢出。 -
Redis 布隆过滤器去重统计 以
区域ID:窗口起始时间作为 Redis 布隆过滤器 Key,隔离不同窗口、不同区域的去重标记;
- 每条车牌通过工具类生成两组哈希偏移位,查询 Redis
getbit判断是否已抓拍过该车; - 若车牌不存在:统计总数 + 1,并用
setbit把对应 bit 位置 1 做已存在标记; - 若车牌已存在:直接跳过,不重复累加;
- 用算子内部
HashMap维护当前【区域 + 窗口】维度的实时去重车辆总数。
本方案优缺点总结
✅ 优点
- 自定义触发器逐条触发计算,窗口无数据积压,彻底解决迭代器大数据量 OOM 问题;
- 基于 Redis 布隆过滤器去重,存储空间极低,适合卡口亿级车流量场景;
- RedisSink 替代 MySQL,大幅提升结果写入吞吐量,消除频繁 JDBC 连接开销;
- 窗口数据生命周期闭环,每条数据即时处理即时释放,内存占用长期平稳可控,满足生产环境长期运行;
- 不同窗口布隆 Key 互相隔离,不会出现跨窗口统计错乱。
❌ 缺点
- 布隆过滤器存在极低概率误判,不能用于要求绝对精准不能少统计的业务;
- 每条数据都要访问 Redis,存在少量网络 IO 损耗,超高并发可部署 Redis 集群优化;
- 本地 HashMap 属于算子内存变量,任务重启后统计数值丢失,如需断点续算需要改成 Flink 状态托管。
六、三套方案整体收尾小结
-
HashSet 窗口批量去重:入门易懂,存在严重内存溢出问题,仅适合测试小流量;
-
Redis 布隆过滤器批量去重:解决存储占用问题,但迭代器数据积压隐患依旧存在;
-
布隆过滤器 + 自定义触发器实时去重:既优化存储开销,又解决数据堆积 OOM,写入性能同步优化,是该业务场景生产最优落地版本。
更多推荐




所有评论(0)