本文目录导读:

这是一个比较专业的分布式系统问题,在 Java 分布式环境中,数据去重的难点在于:数据分散在不同节点、并发写冲突以及跨网络状态同步。
对于“流优化”,通常指对实时数据流(如 Kafka、日志、埋点)进行去重,要求低延迟、高吞吐,且不依赖全量存储。
下面从技术方案和代码实现两个层面,给出几种主流的去重策略及优化思路。
核心思路:去重三要素
无论哪种方案,去重本质上都需要权衡三个要素:存储介质、去重维度 和 时效性。
- 存储介质:内存(快但易失)、Redis(快但有网络IO)、数据库(慢但可靠)。
- 去重维度:唯一ID(如订单号、请求TraceId)、近似的指纹(如MinHash、布隆过滤器)。
- 时效性:精确去重(需要存所有Key)、滑动窗口去重(只存最近N分钟)、概率去重(允许小概率误判)。
主流技术方案及代码示例
方案1:基于 Redis 的 Set / HyperLogLog(最常用)
这是分布式去重的标配,适合中等量级(日活千万级)。
- 精确去重:使用 Redis
Set,sadd key value,利用其唯一性约束。 - 近似去重:使用 Redis
HyperLogLog,pfadd key value,内存极低(12KB可存2^64个元素),误差约0.81%。
优化点:
- Pipeline/Batch:批量提交去重请求,减少网络往返。
- Key 分片:给每个去重 ID 加上时间戳前缀(如
dedup:yyyyMMdd:hash(id)%100),避免单个 Key 过大导致阻塞。 - TTL 过期:设置过期时间(如
expire dedup:yyyyMMdd:hash 86400),自动清理历史数据。
Java 代码示例(Spring Data Redis + Pipeline):
@Service
public class DedupService {
@Autowired
private StringRedisTemplate redisTemplate;
/**
* 批量判断并标记去重
* @param bizKey 业务前缀,如 "order_pay"
* @param ids 待去重的ID列表
* @return 返回新的、不重复的ID列表
*/
public List<String> batchDedup(String bizKey, List<String> ids) {
String today = LocalDate.now().format(DateTimeFormatter.ofPattern("yyyyMMdd"));
List<Object> results = redisTemplate.executePipelined((RedisCallback<Object>) connection -> {
for (String id : ids) {
// 构建key: 业务:日期:分片
String shard = String.valueOf(Math.abs(id.hashCode()) % 100);
String key = "dedup:" + bizKey + ":" + today + ":" + shard;
RedisStringCommands commands = connection.stringCommands();
commands.setNX(key.getBytes(), id.getBytes());
// 设置过期时间,避免无限增长
commands.expire(key.getBytes(), 86400); // 存1天
}
return null;
});
// 处理结果: true表示第一次插入(不重复)
List<String> newIds = new ArrayList<>();
for (int i = 0; i < ids.size(); i++) {
if (Boolean.TRUE.equals(results.get(i))) {
newIds.add(ids.get(i));
}
}
return newIds;
}
}
方案2:基于布隆过滤器(Bloom Filter)—— 极高吞吐,允许误判
适合海量数据(日活亿级)、对绝对精准不敏感的场景(如爬虫去重、已读消息标记)。
- 原理:位数组 + 多个哈希函数,如果所有哈希位置都为1,则可能存在;只要有一个为0,则一定不存在。
- 优化点:使用分层布隆过滤器 + 本地缓存 + Redis 共享。
Java 代码示例(Guava BloomFilter + 本地缓存 + 二级校验):
@Component
public class BloomFilterDedup {
// 本地一级去重: 10万条,误判率1%
private final BloomFilter<String> localFilter = BloomFilter.create(
Funnels.stringFunnel(StandardCharsets.UTF_8), 100_000, 0.01);
// Redis二级共享去重 (用布隆过滤器命令)
@Autowired
private StringRedisTemplate redisTemplate;
public boolean isDuplicate(String id) {
// 1. 先查本地 (极快)
if (localFilter.mightContain(id)) {
// 2. 本地可能存在,再查远端Redis (减少对Redis的QPS)
return isExistsInRedis(id);
} else {
// 3. 本地一定不存在,则去重成功,插入本地和远端
localFilter.put(id);
addToRedis(id);
return false; // 不重复
}
}
private boolean isExistsInRedis(String id) {
// 假设Redis里有一个全局的布隆过滤器 (使用redis-bloom模块的BF.EXISTS命令)
// 这里用伪代码示意
return Boolean.TRUE.equals(
redisTemplate.opsForValue().getBit("bloom:global", getBitIndex(id))
);
// 实际生产推荐用 Redisson 的 RBloomFilter
}
private void addToRedis(String id) {
redisTemplate.opsForValue().setBit("bloom:global", getBitIndex(id), true);
}
private long getBitIndex(String id) {
// 简单哈希取模 (实际应使用多个哈希函数)
return id.hashCode() & 0x7FFFFFFF % 10_000_000;
}
}
方案3:基于数据库(SQL)的唯一索引 + 批量写入(忽略冲突)
适合最终一致性要求不高的场景(如日志、消息记录)。
- 原理:把去重 ID 设为主键或唯一索引,
INSERT IGNORE或ON DUPLICATE KEY UPDATE。 - 优化点:
INSERT IGNORE+ 批量 Batch Insert (JDBC batch)。 - 注意:数据库是性能瓶颈,需要配合分库分表(ShardingSphere/MyCat)。
代码示例(MyBatis + Insert ignore):
<insert id="batchInsertIgnore" parameterType="list">
INSERT IGNORE INTO dedup_log (id, biz_key, create_time) VALUES
<foreach collection="list" item="item" separator=",">
(#{item.id}, #{item.bizKey}, now())
</foreach>
</insert>
// 返回值 = 实际插入行数,小于list.size()的就是重复的
int count = mapper.batchInsertIgnore(dataList);
List<String> dedupedIds = dataList.subList(0, count).stream()
.map(DedupLog::getId).collect(toList());
流式去重优化
如果你是在实时流(如 Flink、Spark Streaming、Kafka Streams)中进行去重,以下是关键优化策略:
状态后端优化(针对 Flink/Spark)
-
不要用全量去重状态,改用 滑动窗口 + BloomFilter。
// Flink 示例:按事件时间开窗,窗口内去重 stream .keyBy(event -> event.getUserId()) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .process(new DedupProcessFunction()); -
使用 RocksDB 状态后端:对于大状态去重(比如一个月内的访问记录),Flink 默认的 HashMap 状态后端会OOM,RocksDB 利用磁盘存储,适合超大状态。
内存去重 + 定时刷盘
- 定时器刷盘:每 1 分钟或每 1000 条,将内存中的去重集合批量写入 Redis/DB。
- Storm/Kafka Streams 里可以用 TimeCacheMap 自动过期。
幂等性设计(终极解法)
如果业务允许,用业务主键 + 幂等消费代替去重。
- Kafka 幂等生产者:
enable.idempotence=true。 - 数据库幂等:写入时
INSERT ... ON DUPLICATE KEY UPDATE。
各种方案对比
| 方案 | 准确度 | 吞吐量 | 内存/存储 | 延迟 | 适合场景 |
|---|---|---|---|---|---|
| Redis Set | 100% | 高(万级/秒) | 高(随数据量增长) | 中等(网络IO) | 中小流量,精确去重 |
| Redis HyperLogLog | 99%+ | 极高(十万级/秒) | 极低(12KB固定) | 中等 | 用户UV统计,模糊去重 |
| Bloom Filter | 95%-99.9% (可调) | 极高(百万级/秒) | 极低(位数组) | 极低(本地) | 爬虫去重、已读标记 |
| 数据库唯一索引 | 100% | 低(千级/秒) | 取决于DB | 高(磁盘IO) | 日志归档、交易记录 |
| Flink/Spark 状态 | 100% (窗口内) | 高(十万级/秒) | 中(RocksDB) | 中等 | 实时流式ETL控重 |
总结建议
先看你的业务容忍度:
- 小规模/低并发,要求精确:直接用 Redis Set + TTL,简单可靠。
- 大规模/高并发,允许极小误差(如爬虫):用 布隆过滤器 做一级过滤,Redis 做二级兜底。
- 实时流控重 / 百万级QPS:用 Flink 窗口状态 + RocksDB + BloomFilter 组合。
- 省钱又持久:用 数据库唯一索引 + 批量模式 + 定期清理(表分区)。
最核心的优化点是:减少跨网络调用(Batch)、缩小状态容量(BloomFilter)、控制生命周期(TTL)。