日均百亿级日志,意味着每秒要处理超过 11 万条数据。在这个量级下,任何内存结构都会瞬间爆掉,任何外部存储都会成为性能瓶颈。
经过多次方案选型和线上实战,我们最终用 Java + RocksDB 扛住了这个挑战。这篇文章,我会把我们的架构设计、踩坑经验和优化实践完整分享出来。
一、百亿级去重的难点在哪?
先说背景:我们的日志系统每天接收上百亿条来自各端的上报数据。由于客户端网络重试、消息队列的 at-least-once 语义等原因,重复消息的比例大约在 0.6% 左右。
对数据准确性要求高的场景来说,这 0.6% 的重复可能意味着数百万的统计误差。
朴素思路很简单:
def dedupe(stream):for message in stream:if message.id in seen_set: # 见过就丢弃discard(message)else:publish(message)seen_set.add(message.id)
问题在于:百亿级 × 1 个月的去重窗口,你根本存不下。
| 存储方案 | 百亿条 ID 的内存占用 | 可行性 |
|---|---|---|
| Java HashSet (UUID, 36字节) | ~360 GB | ❌ 内存爆炸 |
| Redis (String 类型) | ~500 GB + | ❌ 成本太高 |
| 布隆过滤器 | ~12 GB (1%误判率) | ⚠️ 有误差、无法删除 |
我们需要一个精确去重(零容忍误判)、持久化(不丢数据)、低成本(单机百 GB 级别)、高性能(毫秒级读写)的方案。
这些矛盾的需求,最终指向了同一个答案:RocksDB。
二、为什么是 RocksDB?
RocksDB 是 Facebook 开源的嵌入式 KV 存储引擎,基于 LSM-Tree(Log-Structured Merge-Tree)架构。
2.1 它凭什么能扛住百亿级?
1. 写性能极高
LSM-Tree 的核心设计是顺序写。所有写入操作先写内存中的 MemTable,达到阈值后以 SSTable 文件的形式顺序落盘。没有随机写,没有原地更新——这对高吞吐写入场景极其友好。
2. 读取有布隆过滤器加速
每个 SSTable 文件都自带布隆过滤器,可以在访问前快速判断 key 是否可能存在。绝大多数重复消息在内存中被拦截,真正落到磁盘的读请求很少。
3. 天然支持数据淘汰
通过 TTL(Time To Live)机制,可以自动过期超出窗口期的数据,无需手动清理。
4. 嵌入式、零依赖
RocksDB 是一个库(提供 Java JNI 接口),不需要单独部署 Redis 或 MySQL 集群,运维成本极低。
2.2 与主流方案的数据对比
一项公开的性能测试数据很好地说明了差异:
| 去重方案 | 100 万条数据响应时间 |
|---|---|
| MapReduce + HDFS | 194 ms |
| Spark Streaming + MapReduce | 511 ms |
| HBase(全局去重) | 305 ms |
| Flink + RocksDB | 98.7 ms |
RocksDB 方案在百万元素去重场景下,响应时间比其他方案快 2-5 倍。
三、系统架构设计
3.1 整体架构
┌─────────────┐ ┌─────────────┐ ┌─────────────────────────────────┐│ Client │───▶│ Kafka │───▶│ Dedupe Consumer ││ (上报日志) │ │ (输入Topic) │ │ ┌─────────────────────────┐ │└─────────────┘ └─────────────┘ │ │ RocksDB (本地) │ ││ │ - Key: messageId │ ││ │ - Value: timestamp │ ││ │ - TTL: 28天 │ ││ └─────────────────────────┘ │└───────────────┬─────────────────┘│▼┌─────────────────────────────────┐│ Kafka (输出Topic) │└─────────────────────────────────┘
3.2 关键设计决策
1. Kafka 按 messageId 分区
在写入 Kafka 输入 Topic 时,我们按 messageId 进行分区:
// 生产者端:确保相同 messageId 进入同一分区producer.send(new ProducerRecord<>("input-topic",message.getMessageId(), // 分区键message));
这样做的好处是:同一个消费者实例始终处理相同范围的 messageId,RocksDB 的缓存局部性大幅提升。
2. Consumer 与 RocksDB 一一对应
每个 Consumer 实例挂载一个独立的 RocksDB 实例,部署在本地 SSD 上:
public class DedupeConsumer implements Runnable {private RocksDB rocksDB;private KafkaConsumer<String, LogMessage> consumer;public void run() {// 初始化 RocksDBOptions options = new Options();options.setCreateIfMissing(true);options.setKeepLogFileNum(3);// 优化:增大 BlockCacheBlockBasedTableConfig tableConfig = new BlockBasedTableConfig();tableConfig.setBlockCacheSize(512 * 1024 * 1024L); // 512MBoptions.setTableFormatConfig(tableConfig);rocksDB = RocksDB.open(options, "/data/rocksdb/" + partitionId);// 消费循环while (running) {ConsumerRecords<String, LogMessage> records = consumer.poll(1000);for (ConsumerRecord<String, LogMessage> record : records) {processRecord(record);}consumer.commitSync();}}private void processRecord(ConsumerRecord<String, LogMessage> record) {String messageId = record.key();byte[] keyBytes = messageId.getBytes(StandardCharsets.UTF_8);// 查询是否已存在byte[] existing = rocksDB.get(keyBytes);if (existing != null) {// 重复消息,直接丢弃return;}// 新消息:写入 RocksDB 并发送到输出 TopicrocksDB.put(keyBytes,String.valueOf(System.currentTimeMillis()).getBytes());producer.send(new ProducerRecord<>("output-topic", record.value()));}}
3. 使用 TTL 自动淘汰过期数据
百亿级数据不能无限存储。我们设置去重窗口为 28 天,利用 RocksDB 的 TTL 机制自动淘汰:
// 开启 TTL 的 RocksDB 实例public RocksDB openWithTTL(String path, int ttlSeconds) throws RocksDBException {Options options = new Options();options.setCreateIfMissing(true);// 使用 TtlDB 代替普通 RocksDBreturn TtlDB.open(options, path, ttlSeconds, false);}// 使用: 28 天过期RocksDB db = openWithTTL("/data/rocksdb", 28 * 24 * 3600);
四、实战踩坑与优化
4.1 坑一:JNI 内存泄漏
问题:RocksDB 的 Java 接口通过 JNI 调用 C++ 层。每个查询返回的 byte[] 如果不及时处理,会导致 Native 内存堆积。
症状:进程 RSS 正常,但 top 显示 RES 持续增长,最终 OOM Killer。
解决方案:
// 错误:直接使用返回的 byte[]byte[] value = rocksDB.get(keyBytes);String timestamp = new String(value); // value 可能很大// 正确:使用完立即置空,建议 JVM GCbyte[] value = rocksDB.get(keyBytes);if (value != null) {try {processValue(value);} finally {// 帮助 GC 回收value = null;}}
4.2 坑二:布隆过滤器参数不是越大越好
RocksDB 为每个 SSTable 维护布隆过滤器,用于快速判断 key 是否存在。
问题:默认配置下,布隆过滤器可能导致内存爆炸。
优化:根据数据特征调整
BlockBasedTableConfig tableConfig = new BlockBasedTableConfig();// 设置布隆过滤器:10 bits per key(约 1.2% 误判率)tableConfig.setFilterPolicy(new BloomFilter(10, false));// 固定 BlockCache 大小,避免无限增长tableConfig.setBlockCacheSize(512 * 1024 * 1024);// 启用缓存索引和过滤器tableConfig.setCacheIndexAndFilterBlocks(true);
4.3 坑三:Compaction 导致性能抖动
LSM-Tree 在后台进行 Compaction 时,会占用大量 I/O 和 CPU,导致读写延迟飙升。
解决方案:限流 + 错峰
Options options = new Options();// 限制后台 Compaction 线程数options.setMaxBackgroundCompactions(2);// 限制后台 Flush 线程数options.setMaxBackgroundFlushes(1);// 设置写入限流(Rate Limiter)options.setRateLimiter(new RateLimiter(100 * 1024 * 1024)); // 100MB/s
4.4 最终性能数据
经过多轮优化,生产环境的真实表现:
| 指标 | 数值 |
|---|---|
| 日均处理量 | 120 亿条 |
| 单机存储量 | 800 GB(28 天窗口) |
| 平均延迟 | < 5 ms |
| P99 延迟 | < 20 ms |
| 单实例 QPS | 2.5 万+ |
| 准确率 | 100%(无漏去重) |
五、技术选型的启示
回过头看,为什么最终是 Java + RocksDB 扛住了这个场景?
RocksDB 解决的是百亿级状态存储的问题。传统的 HashSet 在千万级就会 OOM,Redis 在百亿级需要上百台机器。而 RocksDB 将数据放在磁盘上,用精巧的 LSM-Tree 结构平衡了读写性能,用布隆过滤器优化了读路径。在最坏情况下,128GB 内存 + 1TB SSD 的单机就能支撑百亿级状态。
Java 解决的是工程落地的问题。RocksDB 提供成熟的 JNI 接口,可以和 Kafka、Flink 等 Java 生态无缝集成。同时,JVM 的内存管理、监控工具(JMX、Arthas)让排查问题变得相对可控。
用一句话总结:用数据库的存储能力,达到缓存的访问延迟,这就是 RocksDB 在百亿级场景下的价值。
写在最后
百亿级去重是一个典型的“海量数据状态管理”问题。我们的实践证明,Java + RocksDB 是一条可行、高效、低成本的路径。
对于同样面临这个问题的团队,我的建议是:
- 优先考虑下推:让去重尽量发生在数据源头,避免在传输链路上引入额外的复杂性。
- 选择合适的窗口:去重窗口不是越长越好,需要结合业务容忍度和存储成本来权衡。
- 压测先行:RocksDB 的配置参数非常多,不同机型、不同负载下最佳配置差异很大,一定要用生产数据做压测。
微信赞赏
支付宝扫码领红包
