一、前言
实时数据流中重复数据是十分常见的问题:上游 MQ 重试、网络抖动、业务重发都会产生重复消息,如果不做去重处理,会导致统计指标虚高、入库数据冗余、业务对账出错。 针对 Flink 实时去重需求,本文循序渐进实现三套方案,由浅入深分析优缺点、内存开销、适用场景:
二、业务背景
当前智慧交通需求:按区域 ID 分组,统计 10 分钟滚动事件时间窗口内,各个区域的独立过车数量。同一辆车短时间多次抓拍会产生重复数据,必须对车牌做去重统计,避免车流量统计虚高。
三、基于 HashSet 内存去重(初级实现)
核心实现思路
核心代码
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);
}
});
方案优缺点分析
✅ 优点
❌ 缺点
存在痛点引出后续优化
为解决 HashSet 内存占用过高问题,我们引入空间利用率更高的布隆过滤器,作为第二版优化方案。
四、Redis 布隆过滤器窗口去重(内存优化版)
什么是布隆过滤器(Bloom Filter)
布隆过滤器是空间利用率极高的概率型二进制数据结构,底层本质是一个超长 bit 位数组,数组内元素只有 0、1 两种值。
核心判定特性
- 如果判定元素一定不存在:返回结果百分百准确;
- 如果判定元素可能存在:存在极小概率哈希碰撞造成误判(不存在的元素被误认为已存在);
- 不存储原始数据本身,只通过多个哈希函数映射标记 bit 位,内存占用远小于 HashSet、Redis Set。
基础插入 & 查询逻辑
- 插入元素:对同一个车牌执行多个不同哈希运算,算出多个 bit 下标,把对应位置全部置为 1;
- 判断重复:取出多个哈希下标,校验所有 bit 位是否全为 1;只要有任意一位是 0,代表车牌从未出现;全部为 1 则判定为已存在。
布隆过滤器与 Redis 的关系
- 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 数据长期堆积。
方案优缺点分析
✅ 优点
❌ 缺点
本方案遗留痛点(引出方案三)
当前 Redis 布隆过滤器绑定固定窗口 Key,虽然解决内存溢出问题,但需要手动管理 Redis 过期清理;如果业务窗口数量极多,Redis Key 会持续累积、运维繁琐。 下一步优化思路:布隆过滤器 + Flink 窗口触发器,在窗口触发结束时自动清理当前窗口布隆过滤器数据,实现生命周期全自动管理,形成生产级稳定去重方案。
五、布隆过滤器 + 自定义触发器 逐条实时去重(生产最终优化方案)
现有前两套方案遗留核心痛点回顾
针对性优化思路: 自定义窗口触发器,每来一条数据立刻触发窗口计算、用完立即清空窗口数据,不让迭代器堆积数据;统计结果落地 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维护当前【区域 + 窗口】维度的实时去重车辆总数。
本方案优缺点总结
✅ 优点
❌ 缺点
六、三套方案整体收尾小结
-
HashSet 窗口批量去重:入门易懂,存在严重内存溢出问题,仅适合测试小流量;
-
Redis 布隆过滤器批量去重:解决存储占用问题,但迭代器数据积压隐患依旧存在;
-
布隆过滤器 + 自定义触发器实时去重:既优化存储开销,又解决数据堆积 OOM,写入性能同步优化,是该业务场景生产最优落地版本。





