拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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

【免费下载链接】flink-learning

flink 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 + TTL

Kafka 数据重复后,各种解决方案思路都比较类似:维护一个全局 Set 集合,存放所有已被处理过的主键。处理新日志时,将当前日志主键与历史 Set 集合比对:

  • 若 Set 中已包含当前主键 → 说明该日志之前已被处理过,过滤掉;
  • 若 Set 中不包含当前主键 → 正常处理,处理完成后将主键加入 Set,使 Set 永远存放所有已被处理过的数据。

处理流程本身很简单,关键在于如何维护这个 Set 集合。以"每天 100 亿数据"为规模基准,可以估算其存储成本:

  1. 主键大小:由于数据量巨大,主键必须足够大以保证不冲突。4 字节 int 只能表示约 42 亿个数,每天百亿数据必然大量冲突,会把不重复数据误判为重复。APP 端没有全局发号器,通常使用 UUID 作为日志主键(36 位字符串,如f106c4a1-4c6f-41c1-9d30-bbb2b271284a),每个主键占 36 字节。
  2. 单日空间:36 字节 × 100 亿 ≈ 360 GB,这仅仅是一天!
  3. 必须加 TTL:若不加 TTL,10 天数据占用 3.6 T,100 天占用 36 T,空间必然爆炸。假设按天去重、重复上报时间间隔不会超过 24 小时,则 TTL 可设为36 小时。
  4. 36 小时窗口内的数据量:100 亿 × 1.5 = 150 亿条主键。
  5. 附带时间戳:每条数据带 TTL 意味着必须额外保存时间戳(如 Redis 中一个 key 设了 TTL 却没有时间戳,就无法判断何时清理)。主键 36 字节 + long 时间戳 8 字节 =每条至少 44 字节。
  6. 总空间: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 万,如何快速判断某个整数是否在其中?

  • 方案 A(HashMap):1000 万个 int 约1000 万 × 4 字节 ≈ 40 MB。
  • 方案 B(boolean 数组):申请长度 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形式存储:

  1. HBase 表设置TTL 为 36 小时,最近 36 小时的 150 亿条日志主键全部以 RowKey 形式存放;
  2. 每来一条数据,先拿主键去 HBase Get 查询:
    • 存在 → 已处理过,过滤;
    • 不存在 → 正常处理,处理完成后将主键 Put 进 HBase 表;
  3. 由于数据量巨大,必须提前对 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 外部 Set,Flink 状态方案的核心优势:

  • 无需额外存储组件:Set 直接作为 Flink 算子状态存在,省去 Redis/HBase 集群的部署与成本;
  • 天然分布式分摊:KeyedState 按 key 分布在各个并行子任务上,百亿主键的存储压力被 TaskManager 集群分摊;
  • 借助 RocksDB 状态后端落盘:使用 RocksDBStateBackend 时,状态存储在本地磁盘(配合内存缓存),成本远低于纯内存,并支持增量 Checkpoint,快照成本可控;
  • 原生 TTL 支持:Flink 的StateTtlConfig可以直接给状态设置 36 小时过期时间,过期数据自动清理,正好对应第 2 节"必须加 TTL"的核心设计;
  • Checkpoint 容错:状态随 Checkpoint 持久化,作业故障恢复后去重记录不丢失。

5.2 如何使用 KeyedState 维护 Set 集合

整体思路:对日志主键做 keyBy,每个 key 对应一个ValueState<Boolean>标记"该主键是否已处理过"。第一次出现时状态为 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并开启增量 Checkpoint;setNumberOfTransferingThreads(3)控制状态恢复/快照时的线程数;预定义模式选择机械硬盘+高内存的组合(SPINNING_DISK_OPTIMIZED_HIGH_MEM)。

② Kafka Consumer 与主键 keyBy:

Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, DeduplicationExampleUtil.broker_list); props.put(ConsumerConfig.GROUP_ID_CONFIG, "keyed-state-deduplication"); FlinkKafkaConsumerBase<String> kafkaConsumer = new FlinkKafkaConsumer<>( DeduplicationExampleUtil.topic, new SimpleStringSchema(), props) .setStartFromGroupOffsets(); env.addSource(kafkaConsumer) .map(log -> GsonUtil.fromJson(log, UserVisitWebEvent.class)) // 反序列化 JSON .keyBy((KeySelector<UserVisitWebEvent, String>) UserVisitWebEvent::getId) .addSink(new KeyedStateSink());

即以日志主键id(UUID 字符串)作为 key 进行keyBy。

③ KeyedStateSink:ValueState + TTL 36 小时(去重核心逻辑):

public static class KeyedStateSink extends RichSinkFunction<UserVisitWebEvent> { // 使用该 ValueState 来标识当前 Key 是否之前存在过 private ValueState<Boolean> isExist; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); ValueStateDescriptor<Boolean> keyedStateDuplicated = new ValueStateDescriptor<>("KeyedStateDeduplication", TypeInformation.of(new TypeHint<Boolean>() { })); // 状态 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 之前已被处理过,当前数据应被过滤 } } }

几个关键设计:

  • 状态类型:ValueState<Boolean>每个 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 ProducerRecord<String, String>(topic, null, null, GsonUtil.toJson(userVisitWebEvent)); producer.send(record);

main方法中循环Thread.sleep(100)+writeToKafka(),即每 100ms 向 Kafka topicuser-visit-log-topic写入一批 JSON 日志。Flink 任务消费同一 topic,Kafka broker 地址与 topic 需与 KeyedStateDeduplication.java 中保持一致。

5.4 优化主键:hash 成 long,减少状态大小并提高吞吐量

keyBy使用 36 位 UUID 字符串,key 本身较长。仓库提供了优化版 TuningKeyedStateDeduplication.java,思路是先用 murmur3_128 把 UUID 字符串 hash 成 long,再以 long 作为 keyBy 的 key:

env.addSource(kafkaConsumer) .map(string -> GsonUtil.fromJson(string, UserVisitWebEvent.class)) // 反序列化 JSON // 这里将日志的主键 id 通过 murmur3_128 hash 后,将生成 long 类型数据当做 key .keyBy((KeySelector<UserVisitWebEvent, 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 相关配置主要有两处:

  1. 预定义优化模式:setPredefinedOptions(PredefinedOptions.SPINNING_DISK_OPTIMIZED_HIGH_MEM)。SPINNING_DISK_OPTIMIZED_HIGH_MEM是"机械硬盘 + 高内存"组合——增加 block cache / write buffer 等内存占用以换取磁盘 IO 减少,适合状态读写密集的去重场景;若集群配备 SSD 且内存有限,可评估换用 SSD 类预定义选项。
  2. 状态迁移线程数:setNumberOfTransferingThreads(3),控制 Checkpoint 上传/恢复时并行传输状态文件的线程数,适当调大可缩短 Checkpoint 时间。
  3. TTL 清理策略:StateTtlConfig中的cleanupInRocksdbCompactFilter(50000000L),在 RocksDB compaction 阶段异步清理过期状态,避免 36 小时 TTL 的过期主键长期占用磁盘(相关配置位于 KeyedStateDeduplication.java)。

调优建议:吞吐量波动时,优先检查 RocksDB 所在磁盘 IO(是否机械盘、是否与 HDFS 写路径争抢 IO)、Checkpoint 是否拖慢主流程(观察minPauseBetweenCheckpoints与实际完成时间),再逐步调整 block cache、write buffer 等 RocksDB 内存参数。

7 小结与反思

本文围绕百亿数据去重,依次对比了四条技术路线:

方案存储成本(百亿/天量级)准确性适用场景
通用全局 Set(Redis 等)≥ 660 GB(内存,成本极高)精准数据量小或预算充足
BloomFilter约 14 GB(10 亿容量、千分之一误判率)有误差可容忍误差的统计场景
HBase 全局 Set660 GB 级廉价存储 + 集群运维精准已有 HBase 集群、接受链路延迟
Flink KeyedState + RocksDB状态落盘、成本低精准流计算内部去重,推荐优先考虑

核心设计原则(贯穿全文):

  1. 去重必须带 TTL:36 小时 TTL 使 Set 规模稳定在 150 亿条量级,否则空间必然爆炸;
  2. 按精度需求选型:能接受误差选 BloomFilter,追求精确选 KeyedState(或 HBase);
  3. 善用 Flink 原生能力:KeyedState +StateTtlConfig+ RocksDBStateBackend(增量 Checkpoint)+ 预定义优化选项,即可在流内部低成本完成百亿级精准去重,无需引入外部存储;
  4. 性能调优从存储介质开始: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_list,topic 为user-visit-log-topic),先运行DeduplicationExampleUtil.main持续造数,再启动KeyedStateDeduplication.main(或TuningKeyedStateDeduplication.main)观察去重效果即可。

  • 示例工程
  • 大数据

【免费下载链接】flink-learning

flink 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
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

返回列表