欢迎光临
我们一直在努力

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)
    未经允许不得转载:171主机测评 » Flink 实时去重三部曲:HashSet → 布隆过滤器 → 布隆 + 窗口触发器渐进式优化实战
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址