资讯动态

基于 Flink 的百亿数据去重实战:从全局 Set、BloomFilter 到 KeyedState 与 RocksDB 优化

发布时间:2026/10/5 1:56:23 来源:尧图企业网站定制
示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载本文内容来自 flink-learning 仓库《Flink 实战与性能优化》第 12.2 节《基于 Flink 的百亿数据去重实践》。文章围绕真实生产环境中Kafka 数据重复这一高频痛点先给出通用去重方案并量化其存储成本再对比 BloomFilter、HBase 全局 Set、Flink KeyedState 三条落地路径最后结合仓库中 flink-learning-project-deduplication 模块的完整源码讲解 KeyedState 去重的实现细节与 RocksDBStateBackend 的调优手段。读完本文你将掌握百亿级去重的空间成本估算方法、BloomFilter 原理与适用边界、Flink 状态去重的完整编码范式以及吞吐量调优的关键参数。1 背景Kafka 中为什么会出现重复数据典型的 APP 用户行为日志分析链路为手机 APP 端 → Nginx 服务端 → Logstash/Flume 等 → Kafka → Flink 任务由于用户手机客户端的网络可能出现不稳定APP 端上传日志普遍遵循宁可重复上报也不能漏报日志的策略因此 Kafka 中同一条日志可能出现 2 条或 2 条以上。而 Flink 任务的数据源基本都是 Kafka当 Kafka 中存在重复数据时实时 ETL 或流计算都必须基于日志主键去重否则会导致计算结果偏高。例如用户 a 在某个页面只点击了一次但点击日志在 Kafka 中出现 2 次最终统计该页面点击数时结果就会偏高。需要注意这里只列举了一种造成重复的可能生产环境中导致 Kafka 数据重复的因素还有很多例如消费端重试、上游重发、手动补数等。本节的核心问题是数据已经重复了该如何处理。仓库中与本主题对应的实战模块是 flink-learning-project-deduplication它提供了完整的可运行去重案例下文将围绕它展开。2 去重的通用解决方案全局 Set TTLKafka 数据重复后各种解决方案思路都比较类似维护一个全局 Set 集合存放所有已被处理过的主键。处理新日志时将当前日志主键与历史 Set 集合比对若 Set 中已包含当前主键 → 说明该日志之前已被处理过过滤掉若 Set 中不包含当前主键 → 正常处理处理完成后将主键加入 Set使 Set 永远存放所有已被处理过的数据。处理流程本身很简单关键在于如何维护这个 Set 集合。以每天 100 亿数据为规模基准可以估算其存储成本主键大小由于数据量巨大主键必须足够大以保证不冲突。4 字节 int 只能表示约 42 亿个数每天百亿数据必然大量冲突会把不重复数据误判为重复。APP 端没有全局发号器通常使用 UUID 作为日志主键36 位字符串如f106c4a1-4c6f-41c1-9d30-bbb2b271284a每个主键占 36 字节。单日空间36 字节 × 100 亿 ≈ 360 GB这仅仅是一天必须加 TTL若不加 TTL10 天数据占用 3.6 T100 天占用 36 T空间必然爆炸。假设按天去重、重复上报时间间隔不会超过 24 小时则 TTL 可设为36 小时。36 小时窗口内的数据量100 亿 × 1.5 150 亿条主键。附带时间戳每条数据带 TTL 意味着必须额外保存时间戳如 Redis 中一个 key 设了 TTL 却没有时间戳就无法判断何时清理。主键 36 字节 long 时间戳 8 字节 每条至少 44 字节。总空间150 亿 × 44 字节 660 GB。结论每天百亿数据量下纯 Set 方案至少需要 660 GB 以上存储空间。这一量化结果决定了后续方案选型的走向要么用低成本的近似结构BloomFilter要么把 Set 放到廉价的分布式存储HBase要么借助 Flink 自身的分布式状态KeyedState RocksDB天然分摊存储。3 方案一BloomFilter 实现去重有些流计算场景对准确性要求并不高例如传统 Lambda 架构中会有离线结果矫正实时结果。当业务可以接受小量误差时可以使用低成本数据结构。BloomFilter 与 HyperLogLog是两类典型代表且都存在误差HyperLogLog 只能估算插入了多少个不重复元素不能回答是否插入过某个元素BloomFilter 恰好相反它能告诉你肯定不包含元素 a或可能包含元素 b但不能告诉你里面插入了多少个元素。3.1 从 bitmap 位图说起假设有 1 千万个整数数据范围 0 ~ 2000 万如何快速判断某个整数是否在其中方案 AHashMap1000 万个 int 约1000 万 × 4 字节 ≈ 40 MB。方案 Bboolean 数组申请长度 2000 万的 boolean 数组以整数为下标值为 true 表示存在。查询整数 K 时直接取array[K]。Java 的 boolean 占 1 字节需要 2000 万字节。方案 C二进制位图用二进制位模拟布尔类型1 表示 true、0 表示 false只需 2000 万 bit ≈2.4 MB是 boolean 数组的 1/8相比 40 MB 原始数据方案更是大幅缩减。但如果数据范围扩大到 0 ~ 100 亿位图就需要 100 亿 bit ≈ 1200 MB反而比存原始数据还大。此时可以压缩位图长度 hash 映射只申请 1 亿 bit对 1000 万个数求 hash 映射到位上约 12 MB但会引入hash 冲突。例如 3 与 100000003 对 1 亿求余都为 3两者映射到同一个 bit如果集合中包含 100000003位图下标 3 为 1查询 3 时就会误判存在。为了减少冲突诞生了 BloomFilter。3.2 BloomFilter 原理hash 冲突无法完全避免于是 BloomFilter 引入多个 hash 函数只要任意一个hash 函数发现元素不在集合中则该元素肯定不在集合中只有当所有hash 函数都命中了才认为元素可能存在于集合中。插入过程以 3 个 hash 函数为例插入元素 a 时3 个函数分别算出位下标 2、8、10将这 3 个位置 1插入元素 b 时算出 5、10、14 并置 1下标 10 被 a、b 共同涉及。查找过程查元素 c 时3 个函数算出 2、6、9其中 6、9 为 0说明 c肯定不在查元素 d 时算出 5、8、14三者均为 1于是判定 d可能存在但实际上 d 从未插入——这正是因为 a、b 恰好覆盖了这 3 个位产生了误判。由此得到 BloomFilter 的两个核心性质不包含的结论是绝对可靠的包含的结论存在一定误判率。误判率与位数组长度、hash 函数个数、已插入元素数量相关插入元素越少误判率越低。3.3 使用 Redis BloomFilter 去重Redis 4.0 之后 BloomFilter 以插件形式加入 Redis。创建时支持设定预期容量预计插入的数据量与误判率插入量达到预期容量时的误判概率。经笔者测试申请预期容量 10 亿、误判率千分之一的 BloomFilter约需143 亿 bit ≈ 14 GB相比 660 GB 的精准 Set 方案存储成本大幅下降。使用中需注意记录 BloomFilter 中已插入的元素个数当插入量达到预期容量10 亿时为了保障误判率应清除当前 BloomFilter 并重新申请一个新的。适用边界BloomFilter 有误差只能用于能承受一定误差的场景如日志分析、UV 类近似统计对于广告计费等对数据精度要求极高的场景应使用精准去重方案。4 方案二HBase 维护全局 Set 实现去重回顾第 2 节的估算百亿精准去重需要维护 150 亿条主键的 Set每条 44 字节共需660 GB 存储空间。注意这里说的是存储空间而非内存空间——660 G 内存太贵660 G 的 Redis 云服务一个月至少 2 万 RMB 以上设计架构必须考虑成本。HBase 基于 RowKey 的 Get 效率很高因此可以将这个大 Set 以HBase RowKey形式存储HBase 表设置TTL 为 36 小时最近 36 小时的 150 亿条日志主键全部以 RowKey 形式存放每来一条数据先拿主键去 HBase Get 查询存在 → 已处理过过滤不存在 → 正常处理处理完成后将主键 Put 进 HBase 表由于数据量巨大必须提前对 HBase 表做预分区将读写压力分散到各个 RegionServer。4.1 HBase RowKey 去重带来的问题从工程实践角度分析该方案主要面临以下挑战本节在原文基础上补充说明预分区策略要求高RowKey 是 UUID 类随机字符串分布虽均匀但预分区数量需按 150 亿量级估算分区过少会造成单个 Region 数据膨胀与热点分区过多则管理成本上升读放大与写放大每条日志至少一次 Get 一次 Put百亿规模下对 HBase 集群的 RPC 压力、WAL 写入压力都很大需要足够的 RegionServer 与合理的表/列族参数支撑成本与运维虽然 HBase 存储比同容量 Redis 内存便宜很多但仍需单独维护一套 HBase 集群或复用已有集群且去重逻辑依赖外部存储链路变长、延迟增加。因此是否选择 HBase 方案需要在额外维护分布式存储与换用 Flink 自身状态之间权衡——后者的优势恰恰在于不需要引入任何外部存储。5 方案三核心使用 Flink KeyedState 实现去重仓库中该方案对应 KeyedStateDeduplication.java是本节重点。5.1 使用 Flink 状态维护 Set 的优势相比 Redis/HBase 外部 SetFlink 状态方案的核心优势无需额外存储组件Set 直接作为 Flink 算子状态存在省去 Redis/HBase 集群的部署与成本天然分布式分摊KeyedState 按 key 分布在各个并行子任务上百亿主键的存储压力被 TaskManager 集群分摊借助 RocksDB 状态后端落盘使用 RocksDBStateBackend 时状态存储在本地磁盘配合内存缓存成本远低于纯内存并支持增量 Checkpoint快照成本可控原生 TTL 支持Flink 的StateTtlConfig可以直接给状态设置 36 小时过期时间过期数据自动清理正好对应第 2 节必须加 TTL的核心设计Checkpoint 容错状态随 Checkpoint 持久化作业故障恢复后去重记录不丢失。5.2 如何使用 KeyedState 维护 Set 集合整体思路对日志主键做 keyBy每个 key 对应一个ValueStateBoolean标记该主键是否已处理过。第一次出现时状态为 null处理后update(true)再次出现时状态非 null直接过滤。完整代码见 KeyedStateDeduplication.java拆解如下。① 环境与状态后端配置主方法前段StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(6); // 使用 RocksDBStateBackend 做为状态后端并开启增量 Checkpoint RocksDBStateBackend rocksDBStateBackend new RocksDBStateBackend( hdfs:///flink/checkpoints, true); rocksDBStateBackend.setNumberOfTransferingThreads(3); // 设置为机械硬盘内存模式强烈建议为 RocksDB 配备 SSD rocksDBStateBackend.setPredefinedOptions( PredefinedOptions.SPINNING_DISK_OPTIMIZED_HIGH_MEM); env.setStateBackend(rocksDBStateBackend); // Checkpoint 间隔为 10 分钟 env.enableCheckpointing(TimeUnit.MINUTES.toMillis(10)); // 配置 Checkpoint CheckpointConfig checkpointConf env.getCheckpointConfig(); checkpointConf.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); checkpointConf.setMinPauseBetweenCheckpoints(TimeUnit.MINUTES.toMillis(8)); checkpointConf.setCheckpointTimeout(TimeUnit.MINUTES.toMillis(20)); checkpointConf.enableExternalizedCheckpoints( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);要点状态后端指向hdfs:///flink/checkpoints并开启增量 CheckpointsetNumberOfTransferingThreads(3)控制状态恢复/快照时的线程数预定义模式选择机械硬盘高内存的组合SPINNING_DISK_OPTIMIZED_HIGH_MEM。② Kafka Consumer 与主键 keyByProperties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, DeduplicationExampleUtil.broker_list); props.put(ConsumerConfig.GROUP_ID_CONFIG, keyed-state-deduplication); FlinkKafkaConsumerBaseString kafkaConsumer new FlinkKafkaConsumer( DeduplicationExampleUtil.topic, new SimpleStringSchema(), props) .setStartFromGroupOffsets(); env.addSource(kafkaConsumer) .map(log - GsonUtil.fromJson(log, UserVisitWebEvent.class)) // 反序列化 JSON .keyBy((KeySelectorUserVisitWebEvent, String) UserVisitWebEvent::getId) .addSink(new KeyedStateSink());即以日志主键idUUID 字符串作为 key 进行keyBy。③ KeyedStateSinkValueState TTL 36 小时去重核心逻辑public static class KeyedStateSink extends RichSinkFunctionUserVisitWebEvent { // 使用该 ValueState 来标识当前 Key 是否之前存在过 private ValueStateBoolean isExist; Override public void open(Configuration parameters) throws Exception { super.open(parameters); ValueStateDescriptorBoolean keyedStateDuplicated new ValueStateDescriptor(KeyedStateDeduplication, TypeInformation.of(new TypeHintBoolean() { })); // 状态 TTL 相关配置过期时间设定为 36 小时 StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.hours(36)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility( StateTtlConfig.StateVisibility.NeverReturnExpired) .cleanupInRocksdbCompactFilter(50000000L) .build(); // 开启 TTL keyedStateDuplicated.enableTimeToLive(ttlConfig); // 从状态后端恢复状态 isExist getRuntimeContext().getState(keyedStateDuplicated); } Override public void invoke(UserVisitWebEvent value, Context context) throws Exception { // 当前 key 第一次出现时isExist.value() 会返回 null if (null isExist.value()) { // ... 这里执行代码处理的逻辑 // 执行完处理逻辑后更新状态值 isExist.update(true); } else { // isExist.value() 不为 null表示当前 key 之前已被处理过当前数据应被过滤 } } }几个关键设计状态类型ValueStateBoolean每个 key 只存 1 个布尔值配合 RocksDB 落盘空间成本极低远小于 44 字节/条的通用 Set 估算TTL 配置过期时间 36 小时UpdateType.OnCreateAndWrite表示只在创建和写入时更新过期时间读取不刷新避免持续活跃的主键永不失效StateVisibility.NeverReturnExpired保证已过期状态不会被读到cleanupInRocksdbCompactFilter(50000000L)表示在 RocksDB compaction 时对过期状态进行清理每处理 5000 万条状态条目触发一次避免 TTL 过期数据长期占用磁盘判断逻辑isExist.value() null表示 key 第一次出现 → 执行正常处理逻辑并update(true)非 null 表示重复 → 过滤。5.3 数据模型与模拟数据生成去重的日志模型为 UserVisitWebEvent.java包含id日志唯一 id、date日期如 20191025、pageId页面 id、userId用户唯一标识、url页面 url五个字段。模拟数据生成器为 DeduplicationExampleUtil.java核心逻辑public static final String broker_list 192.168.30.215:9092,192.168.30.216:9092,192.168.30.220:9092; public static final String topic user-visit-log-topic; // 每个事件用 UUID 生成唯一 id其余字段随机生成 UserVisitWebEvent userVisitWebEvent UserVisitWebEvent.builder() .id(UUID.randomUUID().toString()) // 日志的唯一 id .date(yyyyMMdd) // 日期 .pageId(pageId) // 页面 id .userId(Integer.toString(userId)) // 用户 id .url(url/ pageId) // 页面的 url .build(); ProducerRecord record new ProducerRecordString, String(topic, null, null, GsonUtil.toJson(userVisitWebEvent)); producer.send(record);main方法中循环Thread.sleep(100)writeToKafka()即每 100ms 向 Kafka topicuser-visit-log-topic写入一批 JSON 日志。Flink 任务消费同一 topicKafka broker 地址与 topic 需与 KeyedStateDeduplication.java 中保持一致。5.4 优化主键hash 成 long减少状态大小并提高吞吐量keyBy使用 36 位 UUID 字符串key 本身较长。仓库提供了优化版 TuningKeyedStateDeduplication.java思路是先用 murmur3_128 把 UUID 字符串 hash 成 long再以 long 作为 keyBy 的 keyenv.addSource(kafkaConsumer) .map(string - GsonUtil.fromJson(string, UserVisitWebEvent.class)) // 反序列化 JSON // 这里将日志的主键 id 通过 murmur3_128 hash 后将生成 long 类型数据当做 key .keyBy((KeySelectorUserVisitWebEvent, Long) log - Hashing.murmur3_128(5).hashUnencodedChars(log.getId()).asLong()) .addSink(new KeyedStateDeduplication.KeyedStateSink());该优化带来的收益可以从实现推断减少 key 序列化与比较成本long 固定 8 字节远小于 36 字符的 UUID 字符串key 的序列化开销、网络 shuffle 开销与状态内部比较开销都更小有利于提升吞吐量key 分布更均匀murmur3_128 是高质量哈希能把 UUID 字符串均匀散列到 long 空间配合并行度 6 可使各并行子任务负载更均衡去重语义不变对相同字符串 hash 结果恒定同一条日志的主键仍会落到同一个 key 上ValueState去重逻辑完全复用KeyedStateDeduplication.KeyedStateSink。注意hash 理论上存在极小概率的碰撞此优化适用于可容忍该概率的业务与 BloomFilter 的取舍类似对精度要求严格的场景建议直接用字符串主键。另外优化版通过.setStartFromLatest()从最新 offset 消费便于联调观察。6 使用 RocksDBStateBackend 的优化方法运行上述方案时如果出现吞吐量时高时低、或实测吞吐量偏低的情况可以从以下几个方面调优对应原文 12.2.5 节以下结合仓库源码展开。6.1 设置本地 RocksDB 的数据目录RocksDBStateBackend 把状态存到 TaskManager 本地磁盘RocksDB 的本地数据目录可通过 Flink 配置项state.backend.rocksdb.local-dir指定。百亿去重场景状态量大、读写频繁建议将本地目录挂载到SSD源码注释也明确建议强烈建议为 RocksDB 配备 SSD配置多块磁盘目录逗号分隔让 RocksDB 在多块盘间分摊 IO缓解单盘瓶颈避免与系统盘、日志盘共用目录防止 IO 互相干扰。6.2 Checkpoint 参数相关配置KeyedStateDeduplication.java 中给出了完整可用的 Checkpoint 配置参数含义如下参数代码配置说明Checkpoint 间隔enableCheckpointing(10 min)每隔 10 分钟触发一次 Checkpoint间隔越大快照频率越低、对吞吐影响越小但故障恢复丢失的数据窗口越大Checkpoint 模式EXACTLY_ONCE精确一次语义配合去重保证结果不重不丢两次 Checkpoint 最小间隔setMinPauseBetweenCheckpoints(8 min)防止频繁触发 Checkpoint给业务处理留出稳定时间片Checkpoint 超时setCheckpointTimeout(20 min)超过 20 分钟未完成则视为失败防止快照卡住拖垮作业外部 CheckpointRETAIN_ON_CANCELLATION作业取消后保留 Checkpoint便于手动恢复或迁移增量 Checkpoint 已通过new RocksDBStateBackend(checkpointPath, true)的第二个参数开启只上传新增/变更的 state 文件大幅减少百亿状态下的快照传输量。6.3 RocksDB 参数相关配置仓库代码中的 RocksDB 相关配置主要有两处预定义优化模式setPredefinedOptions(PredefinedOptions.SPINNING_DISK_OPTIMIZED_HIGH_MEM)。SPINNING_DISK_OPTIMIZED_HIGH_MEM是机械硬盘 高内存组合——增加 block cache / write buffer 等内存占用以换取磁盘 IO 减少适合状态读写密集的去重场景若集群配备 SSD 且内存有限可评估换用 SSD 类预定义选项。状态迁移线程数setNumberOfTransferingThreads(3)控制 Checkpoint 上传/恢复时并行传输状态文件的线程数适当调大可缩短 Checkpoint 时间。TTL 清理策略StateTtlConfig中的cleanupInRocksdbCompactFilter(50000000L)在 RocksDB compaction 阶段异步清理过期状态避免 36 小时 TTL 的过期主键长期占用磁盘相关配置位于 KeyedStateDeduplication.java。调优建议吞吐量波动时优先检查 RocksDB 所在磁盘 IO是否机械盘、是否与 HDFS 写路径争抢 IO、Checkpoint 是否拖慢主流程观察minPauseBetweenCheckpoints与实际完成时间再逐步调整 block cache、write buffer 等 RocksDB 内存参数。7 小结与反思本文围绕百亿数据去重依次对比了四条技术路线方案存储成本百亿/天量级准确性适用场景通用全局 SetRedis 等≥ 660 GB内存成本极高精准数据量小或预算充足BloomFilter约 14 GB10 亿容量、千分之一误判率有误差可容忍误差的统计场景HBase 全局 Set660 GB 级廉价存储 集群运维精准已有 HBase 集群、接受链路延迟Flink KeyedState RocksDB状态落盘、成本低精准流计算内部去重推荐优先考虑核心设计原则贯穿全文去重必须带 TTL36 小时 TTL 使 Set 规模稳定在 150 亿条量级否则空间必然爆炸按精度需求选型能接受误差选 BloomFilter追求精确选 KeyedState或 HBase善用 Flink 原生能力KeyedState StateTtlConfig RocksDBStateBackend增量 Checkpoint 预定义优化选项即可在流内部低成本完成百亿级精准去重无需引入外部存储性能调优从存储介质开始RocksDB 强烈建议配 SSD同时关注 Checkpoint 频率与 RocksDB 内存/IO 参数。如需进一步阅读与运行文章原文books/flink-in-action-12.2.md《Flink 实战与性能优化》第 12.2 节去重实战模块flink-learning-project-deduplication精准去重实现KeyedStateDeduplication.java主键 hash 优化实现TuningKeyedStateDeduplication.java日志模型与模拟数据UserVisitWebEvent.java、DeduplicationExampleUtil.java运行方式将模拟数据生成器与 Flink 任务部署到可访问 Kafka 的环境broker 地址见DeduplicationExampleUtil.broker_listtopic 为user-visit-log-topic先运行DeduplicationExampleUtil.main持续造数再启动KeyedStateDeduplication.main或TuningKeyedStateDeduplication.main观察去重效果即可。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐基于 Apache Flink 的百亿数据去重实践KeyedState TTL 与 RocksDB 调优flink-learning 项目实战基于 Apache Flink 的百亿数据去重实践KeyedState TTL 与 RocksDB 调优flink learning 项目实战 本篇文示例工程大数据百亿数据实时去重Flink状态后端选型与优化百亿数据实时去重Flink状态后端选型与优化 本文深入探讨了在百亿级别数据量下实现实时去重的技术挑战与解决方案。文章首先分析了用户行为日志分析、电商交易订单去示例工程大数据Flink SQL 去重Deduplication基于 ROW_NUMBER() 的流式去重写法与优化器实现Flink SQL 去重Deduplication基于 ROW_NUMBER 的流式去重写法与优化器实现 本文以 Flink 官方文档《Deduplica后端大数据流处理批处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑