Java分布式数据去重流优化等怎么去重

wen java案例 20

本文目录导读:

Java分布式数据去重流优化等怎么去重

  1. 核心思路:去重三要素
  2. 主流技术方案及代码示例
  3. 流式去重优化
  4. 各种方案对比
  5. 总结建议

这是一个比较专业的分布式系统问题,在 Java 分布式环境中,数据去重的难点在于:数据分散在不同节点并发写冲突以及跨网络状态同步

对于“流优化”,通常指对实时数据流(如 Kafka、日志、埋点)进行去重,要求低延迟、高吞吐,且不依赖全量存储。

下面从技术方案代码实现两个层面,给出几种主流的去重策略及优化思路。


核心思路:去重三要素

无论哪种方案,去重本质上都需要权衡三个要素:存储介质去重维度时效性

  1. 存储介质:内存(快但易失)、Redis(快但有网络IO)、数据库(慢但可靠)。
  2. 去重维度:唯一ID(如订单号、请求TraceId)、近似的指纹(如MinHash、布隆过滤器)。
  3. 时效性:精确去重(需要存所有Key)、滑动窗口去重(只存最近N分钟)、概率去重(允许小概率误判)。

主流技术方案及代码示例

方案1:基于 Redis 的 Set / HyperLogLog(最常用)

这是分布式去重的标配,适合中等量级(日活千万级)。

  • 精确去重:使用 Redis Setsadd key value,利用其唯一性约束。
  • 近似去重:使用 Redis HyperLogLogpfadd 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 IGNOREON 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控重

总结建议

先看你的业务容忍度:

  1. 小规模/低并发,要求精确:直接用 Redis Set + TTL,简单可靠。
  2. 大规模/高并发,允许极小误差(如爬虫):用 布隆过滤器 做一级过滤,Redis 做二级兜底。
  3. 实时流控重 / 百万级QPS:用 Flink 窗口状态 + RocksDB + BloomFilter 组合。
  4. 省钱又持久:用 数据库唯一索引 + 批量模式 + 定期清理(表分区)。

最核心的优化点是:减少跨网络调用(Batch)、缩小状态容量(BloomFilter)、控制生命周期(TTL)

抱歉,评论功能暂时关闭!