You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何安全实现CQRS模式下基于轮询的新增数据拉取?

解决CQRS模式下消费者安全轮询新增记录的方案

嘿,这个问题在CQRS落地的时候真的挺常见的——毕竟多生产者插数据、多消费者不同速度拉取,很容易碰到漏数据或者重复处理的坑。我结合实际项目经验,给你梳理几个安全可靠的实现思路:

1. 基于严格递增唯一ID的分页轮询

这是最直接也最稳妥的方案,前提是你的记录有严格单调递增的唯一ID(比如MySQL自增主键、分布式场景下的雪花ID/UUIDv7这类带时间戳的递增ID)。

核心逻辑是:每个消费者维护自己上次处理的最大ID,每次轮询时直接查询ID > 上次最大ID的记录,再按ID排序批量拉取。比如SQL可以这么写:

SELECT * FROM event_records 
WHERE id > ? 
ORDER BY id ASC 
LIMIT 100; -- 批量大小根据业务调整

为什么靠谱?因为记录插入后不会修改,递增ID能保证你不会漏过任何一条新插入的记录——只要新记录的ID比上次的大,就一定会被查到。而且多个消费者可以各自维护自己的ID进度,互不干扰。

2. 时间戳+唯一ID的复合查询(解决时间戳精度问题)

如果你的ID不是严格递增的(比如用了随机UUID),或者担心时间戳的精度问题(比如同一毫秒插入多条记录),可以把时间戳和ID组合起来用。

每个消费者维护两个进度值:上次处理的最大时间戳和对应时间戳下的最大ID。查询条件就变成:

SELECT * FROM event_records 
WHERE (created_at > ?) OR (created_at = ? AND id > ?)
ORDER BY created_at ASC, id ASC
LIMIT 100;

这样即使同一时间戳下有多条记录,也能通过ID继续筛选,不会漏掉任何一条,完美解决时间戳精度不够的问题。

3. 用数据库CDC(变更数据捕获)替代轮询

如果轮询的延迟或者效率满足不了你的需求,那直接用数据库的CDC机制更省心。比如:

  • MySQL可以监听Binlog
  • PostgreSQL用Wal2Json或者Logical Replication
  • SQL Server开Change Tracking

CDC会实时捕获数据库的所有插入操作,消费者可以订阅这些变更事件,不需要主动轮询。数据库层面会保证所有插入的记录都能被捕获,完全不用担心漏数据的问题,而且实时性比轮询好太多。

4. 为消费者维护持久化的进度表

不管用哪种轮询方式,都建议把消费者的进度持久化到数据库里,比如专门建一张consumer_progress表:

CREATE TABLE consumer_progress (
    consumer_id VARCHAR(50) PRIMARY KEY,
    last_processed_id BIGINT NOT NULL,
    last_processed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);

每次处理完一批记录后,原子性地更新这个表的进度。比如用带条件的UPDATE语句:

UPDATE consumer_progress 
SET last_processed_id = ?, last_processed_at = NOW()
WHERE consumer_id = ? AND last_processed_id = ?;

这样即使消费者进程重启,也能从上次的进度继续,不会重复处理或者遗漏。多个消费者可以用不同的consumer_id来区分各自的进度。

5. 幂等性兜底:避免重复处理

最后,不管用哪种方案,都建议给每条记录加一个唯一的业务标识(比如event_uuid),同时消费者维护一个已处理记录的标识列表(可以存在本地缓存或者专门的processed_events表)。

处理记录前先检查这个标识是否已经存在,如果存在就直接跳过。这是兜底方案,防止因为网络波动、进程重启等意外情况导致的重复消费,保证系统的幂等性。

内容的提问来源于stack exchange,提问作者Gili

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 08:29:55