日均百亿级日志,意味着每秒要处理超过 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       │    │                                       │  │  - Valuetimestamp     │    │                                       │  │  - 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<StringLogMessage> consumer;
    @Override    public void run() {        // 初始化 RocksDB        Options options = new Options();        options.setCreateIfMissing(true);        options.setKeepLogFileNum(3);
        // 优化:增大 BlockCache        BlockBasedTableConfig tableConfig = new BlockBasedTableConfig();        tableConfig.setBlockCacheSize(512 * 1024 * 1024L);  // 512MB        options.setTableFormatConfig(tableConfig);
        rocksDB = RocksDB.open(options, "/data/rocksdb/" + partitionId);
        // 消费循环        while (running) {            ConsumerRecords<StringLogMessage> records = consumer.poll(1000);            for (ConsumerRecord<StringLogMessage> record : records) {                processRecord(record);            }            consumer.commitSync();        }    }
    private void processRecord(ConsumerRecord<StringLogMessage> record) {        String messageId = record.key();        byte[] keyBytes = messageId.getBytes(StandardCharsets.UTF_8);
        // 查询是否已存在        byte[] existing = rocksDB.get(keyBytes);        if (existing != null) {            // 重复消息,直接丢弃            return;        }
        // 新消息:写入 RocksDB 并发送到输出 Topic        rocksDB.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 代替普通 RocksDB    return 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(10false));// 固定 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 是一条可行、高效、低成本的路径。

对于同样面临这个问题的团队,我的建议是:

  1. 优先考虑下推:让去重尽量发生在数据源头,避免在传输链路上引入额外的复杂性。
  2. 选择合适的窗口:去重窗口不是越长越好,需要结合业务容忍度和存储成本来权衡。
  3. 压测先行:RocksDB 的配置参数非常多,不同机型、不同负载下最佳配置差异很大,一定要用生产数据做压测。
扫码领红包

微信赞赏支付宝扫码领红包

发表回复

后才能评论