
在数据同步领域你很可能听说过两种方案Cursor-based sync基于游标的同步和Change Data CaptureCDC变更数据捕获。很多团队在选型时会直接在游标同步和 CDC 之间二选一但实际落地时却发现游标同步好像“缺”了删除事件CDC 又对运维要求不低。这篇文章会讲清楚两个方案的核心原理、适用边界以及你在对比时最容易忽略的三个点删除和 DDL 事件、数据一致性语义、对源库的性能影响。文章不会只停在概念层面。我会先给你一张速览表再分别用 SQL 和 Python 演示游标同步的常见写法接着给出 Debezium / Flink CDC 的通用配置模板最后用一套可重复的测试流程验证两种方案在新增、更新、删除场景下的表现。内容适合正在做数据同步、实时数仓、多活架构的开发和架构师阅读。如果你只是在选型阶段这篇文章也能帮你省掉不少试错时间。1. 核心能力速览对比项游标同步Cursor-based变更数据捕获CDC捕获方式按条件分批读取源表数据记录游标位置解析数据库日志如 binlog、WAL或使用触发器/轮询典型实现WHERE id ? ORDER BY id LIMIT ?Debezium、Flink CDC、Canal、Maxwell、AWS DMS实时性取决于轮询间隔通常秒级到分钟级可做到毫秒级依赖日志落盘和解析速度增量数据新增、更新可识别删除通常需要特殊处理支持插入、更新、删除甚至 DDL 变更对源库影响轮询会产生查询压力特别是大表范围查询日志解析对业务压力小但需要开启 binlog/WAL架构复杂度低依赖数据库连接和任务调度高需要部署 connector / 消息队列 / 状态管理数据一致性基于业务字段可能丢失中间状态一般以数据库事务为边界事件顺序更可靠适用场景离线批同步、简单单向同步、小数据量实时数仓、数据订阅、多副本同步、需捕获删除场景从表格能看到两者不是简单的“新旧替代”关系而是取舍问题。游标同步最容易被团队接受因为不需要额外组件但 CDC 能补上游标同步抓不到的数据变化细节。2. 适用场景与使用边界游标同步适合谁适合数据量不大、对实时性要求不高的离线批处理场景。比如每天同步一次报表库、把订单表同步到搜索引擎、或者做简单的数据清洗。它的实现成本很低只要能在源库执行SELECT就能写同步逻辑排错也直观查看游标位置便知道同步进度。CDC 适合谁适合需要低延迟获取增量数据的场景比如实时数仓、实时风控、跨机房数据同步、审计日志等。因为 CDC 能捕获删除和更新前镜像所以很多强一致性场景必须用它。游标同步的边界恰恰是很多人“没看到”的地方删除事件会丢。如果业务表有物理删除游标同步默认无法发现。更新后无法感知旧值。如果你需要知道“这个字段从上个值改成新值”游标同步只能读到最新值。DDL 变更会断任务。比如增加一列游标同步不会自动处理通常需要人工干预。频繁轮询大表会带来性能压力。特别是没有索引的WHERE条件容易拖垮源库。CDC 的边界也要说清楚要求源数据库开启日志并且日志格式要匹配工具要求。需要维护 connector 和消息队列组件变多后排错链路变长。部分 CDC 工具对数据库用户权限有要求比如需要REPLICATION SLAVE、SELECT权限。日志保留时间内如果消费中断可能需要重新初始化否则丢失事件。合规与授权方面如果你是同步业务数据需要确保有数据库访问授权符合数据安全和个人信息保护的相关要求。如果数据来自用户处理时要做好脱敏和审计不能因为“只是同步”就忽略隐私保护。3. 游标同步Cursor-based sync原理与实现3.1 游标同步的基本模型游标同步的核心是维护一个“游标值”每次从源表读取一批数据处理完后推进游标再读取下一批。这个过程不是数据库CURSOR的概念而是一种分页增量读取的工程惯用法。最简单的游标值有两个选择主键自增 ID更新时间updated_at主键适合只能处理新增无法处理更新更新时间可以处理更新但要求业务表必须维护updated_at并且更新时要保证该字段会被正确写入。如果你需要兼顾新增和更新通常会同时使用“主键 更新时间”的组合条件。3.2 基于自增 ID 的游标同步假设源表结构如下CREATE TABLE user_order ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10,2), status VARCHAR(32), created_at DATETIME );同步时每一批执行的 SQL 是SELECT * FROM user_order WHERE id #{last_id} ORDER BY id ASC LIMIT #{batch_size};关键点在于last_id必须持久化比如存放在同步任务的元数据表或本地文件里。如果一个事务里插入了多行并且 ID 单调递增这种方案不会丢数据。但如果业务上存在 Redis 预生成 ID或者用分布式 IDID 顺序与业务时间顺序可能不一致这时候增量数据的先后顺序就变得不确定。3.3 基于更新时间的游标同步如果业务表有updated_at可以使用SELECT * FROM user_order WHERE updated_at #{last_sync_time} ORDER BY updated_at ASC LIMIT #{batch_size};这里有几个常见的坑同一条记录在一个批次同步期间被更新了多次如果更新时间相同可能漏记录。有批量 UPDATE 语句不会触发 updated_at 更新除非程序强制设置。与原库时区不一致会导致比较偏差建议统一用 UTC 时间戳。被删除的数据永远不会出现在结果里。为了避免时间相同导致漏数据你需要把last_sync_time当作“截止时间”而不是“起始时间”并且执行时使用同时结合主键去重。更稳妥的做法是让业务表增加一个版本号每次更新都递增游标基于版本号推进。3.4 Python 示例代码下面给出一个通用游标同步脚本模板。它从源表读取增量数据打印每批行数并把新的游标位置写进本地 JSON 文件。实际项目里你可以把处理逻辑替换成写入目标库或 MQ。import json import os import time import pymysql CURSOR_FILE ./sync_cursor.json BATCH_SIZE 100 POLL_INTERVAL 5 # 秒 def load_cursor(): if os.path.exists(CURSOR_FILE): with open(CURSOR_FILE, r) as f: return json.load(f) return {last_id: 0, last_updated: 1970-01-01 00:00:00} def save_cursor(cursor): with open(CURSOR_FILE, w) as f: json.dump(cursor, f) def fetch_orders(conn, last_id, batch_size): with conn.cursor() as cur: sql SELECT id, user_id, amount, status, updated_at FROM user_order WHERE id %s ORDER BY id ASC LIMIT %s cur.execute(sql, (last_id, batch_size)) return cur.fetchall() def main(): conn pymysql.connect( host127.0.0.1, usersync_user, passwordyour_password, databasesource_db, charsetutf8mb4, ) cursor load_cursor() while True: rows fetch_orders(conn, cursor[last_id], BATCH_SIZE) if not rows: print(f[{time.strftime(%H:%M:%S)}] 没有新数据休眠 {POLL_INTERVAL}s) time.sleep(POLL_INTERVAL) continue for row in rows: # 这里写你的业务处理逻辑比如写入目标库 print(处理订单:, row[0], row[1], row[2], row[3]) # 推进游标取本批最大 id cursor[last_id] rows[-1][0] save_cursor(cursor) print(f本批同步 {len(rows)} 条游标推进到 {cursor[last_id]}) if len(rows) BATCH_SIZE: time.sleep(POLL_INTERVAL) if __name__ __main__: main()注意这个示例只处理新增场景。如果要处理更新你需要把updated_at也作为游标条件并且记录最后一条同步时间。另外生产环境不要用本地 JSON 保存游标建议放在元数据库或 Redis 中避免多实例同时运行时互相覆盖。4. Change Data CaptureCDC原理与实现4.1 CDC 的核心机制CDC 全称是 Change Data Capture意思是“捕获数据变更”。它不是一种单一技术而是一类方案。常见的实现方式有基于数据库日志解析例如 MySQL binlog、PostgreSQL WAL基于数据库触发器在业务表上创建触发器把变更写入一张变更表基于轮询快照对比定期比较不同时间点的数据镜像生产环境中日志解析是主流。因为触发器会严重影响数据库写性能而快照对比的实时性太差。日志解析能够以较低成本获得接近实时的数据变更流。以 MySQL binlog 为例打开binlog_formatROW后binlog 里每一条变更事件都会包含完整行的前后镜像。Canal、Debezium、Flink CDC 等工具都能解析这些事件并把它们转换成统一的结构化消息。4.2 基于数据库日志的 CDC使用 CDC 的核心流程是在源数据库开启 binlog或 WAL并设置正确的格式。使用连接器连接到源数据库读取日志。连接器把解析出来的变更记录发送到消息队列或直接推送给消费者。消费者根据事件类型insert、update、delete做业务处理。这里举一个 MySQL 开启 binlog 的示例配置# /etc/mysql/mysql.conf.d/mysqld.cnf [mysqld] server-id 223344 log-bin mysql-bin binlog_format ROW expire_logs_days 7expire_logs_days代表 binlog 保留天数具体数值要根据业务量和消费速度设置。如果同步任务挂了几天binlog 已经被清理重新启动后可能要从头开始或者需要重新执行全量快照。这是 CDC 任务最常见的一种恢复场景。4.3 Debezium / Flink CDC 的通用配置Debezium 是目前使用最广泛的 CDC 连接器之一。下面是一个 Debezium MySQL connector 的配置模板JSON 格式{ name: mysql-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 127.0.0.1, database.port: 3306, database.user: debezium_user, database.password: your_password, database.server.id: 223345, database.server.name: mysql-server, database.allowPublicKeyRetrieval: true, database.history.kafka.bootstrap.servers: 127.0.0.1:9092, database.history.kafka.topic: schema-changes.mysql-server, table.include.list: source_db.user_order, include.schema.changes: true, snapshot.mode: initial } }如果你的技术栈是 Java 和 Flink那么更常用的是 Flink CDC。Flink CDC 可以把数据库变更直接转换成 DataStream下面是一个极简的 Flink SQL 示例CREATE TABLE user_order ( id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, primary key (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 127.0.0.1, port 3306, username flink_user, password your_password, database-name source_db, table-name user_order ); -- 将变更写入目标表 INSERT INTO sync_target SELECT * FROM user_order;这里并不需要你去记忆具体依赖版本实际项目中要根据 Flink 版本和 CDC 版本确定mysql-cdc的参数。重点是明白CDC 工具会把数据库的 binlog 流包装成标准的事件流你再也不用担心删除事件丢失。5. 两者对比你在错过什么很多团队选型时只会看“哪个接入简单”下意识地选了游标同步。但下面这些差异才是真正影响系统设计的地方。5.1 延迟游标同步的延迟取决于轮询频率。如果每秒轮询一次理论延迟在 1 秒左右但为了控制数据库压力你通常会把轮询间隔放到 10 秒以上延迟自然变成几十秒。如果业务对实时性有秒级以下要求游标同步基本做不到。CDC 的延迟主要受日志解析和传输链路影响。正常情况下binlog 解析到消息中间件的时间在毫秒级到秒级。如果消息队列出现积压延迟会变大但一般仍优于周期性轮询。5.2 数据一致性游标同步的“一致性”很难保障。举例来说一条记录在一个批次内被更新了多次游标同步可能只读到最终值中间状态全部丢失。如果你需要分析“从待付款变成已付款”的事件本身或者需要记录变更历史游标同步只能靠全字段快照来猜并不准确。CDC 则能提供更细粒度的事件数据。以update事件为例Debezium 会同时携带before和after字段你可以清楚看到变更前后的值。这正是“你 missing 的部分”——很多同步需求并不仅仅是要最终结果而是需要一条完整的事件流。5.3 对源库的影响游标同步对源库的查询压力不可忽视。批量SELECT对大表来说如果条件没有索引或者涉及范围扫描会占用大量 IO。这在高频轮询时尤其明显。CDC 日志解析对源库的压力小得多但会增加日志的写入量并需要额外的用户权限。如果源库本身的磁盘和 IO 已经紧张开启 binlog 也可能带来新的风险。整体来说CDC 的性能模型更可控。5.4 删除和 DDL 的支持这是最容易被忽略的点。游标同步依赖于SELECT所以物理删除天然无法感知。如果你要同步的目标系统也需要处理“删除”典型的例子是搜索引擎必须删除对应文档游标同步只能靠额外标记比如deleted_at字段或定时全量覆盖来解决。CDC 则能直接输出delete事件。你可以在消费者中收到删除事件后去目标库执行对应删除或者投递到消息队列供下游处理。对于 DDLCDC 工具通常也能捕获到表结构变更方便你维护元数据。5.5 复杂性必须承认CDC 的架构复杂度更高。你需要部署连接器、消息队列、消费者任务还要处理断点续传、Schema 变更、数据反序列化等。游标同步往往一条 SQL 加一个定时程序就能跑起来。所以选择哪一个方案取决于你愿不愿意为“完整变更流”支付运维成本。如果业务只要求最终一致且能接受批延迟游标同步并没有错。但如果你面临实时数仓或需要订阅增量事件强行用游标同步去补删除逻辑最后会发现成本一点不比 CDC 低。6. 实践环境与验证流程我不会给你一个虚构的压测数据但会给你一套可落地的验证流程。你可以用这套流程在自己的测试库上跑一遍观察两种方案的真实表现。6.1 准备数据库测试环境你需要准备一个 MySQL 实例或其他支持日志格式的数据库开启 binlog并将格式设置为ROW创建一张测试表插入一批初始数据安装 Python 和 PyMySQL 用于编写游标同步脚本如果你要测 CDC建议准备 Docker 环境运行 Debezium 或 Flink CDC 容器创建测试表的 SQLCREATE TABLE user_order ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10,2), status VARCHAR(32), updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); INSERT INTO user_order (user_id, amount, status) VALUES (1001, 199.00, PAID), (1002, 299.00, UNPAID), (1003, 399.00, REFUNDED), (1004, 499.00, PAID);6.2 游标同步功能测试运行第 3.4 节的 Python 脚本观察第一批能同步 4 条数据。再向表中插入一条新记录脚本在下个轮询周期会读到第 5 条记录。执行一条UPDATE把status改成CANCELED。如果脚本只依赖id它不会感知这条更新如果你基于updated_at需要额外扩展游标条件。执行一条DELETE观察脚本是否完全不知道删除发生。这个测试就可以直观告诉你游标同步的边界能新增、能更新条件满足时但删除不会出现。6.3 CDC 功能测试使用 Debezium 或 Flink CDC 时建议按下面步骤验证启动 CDC 连接器确认它先做一次 snapshot把当前 4 条数据发给下游。执行INSERT下游能收到一条create事件。执行UPDATE下游能收到一条update事件事件里包含before和after两个字段。执行DELETE下游能收到一条delete事件。如果连接器配置了include.schema.changes true执行ALTER TABLE后你还会收到一个 DDL 事件。判断是否成功的标准是下游消费者收到的消息顺序、字段、事件类型是否与数据库实际操作一致。你可以把消息打印到日志里确认每一条变更都被记录。6.4 效果判断标准测试动作游标同步预期表现CDC 预期表现插入新数据出现在后续批查询收到 insert/create 事件更新依赖游标字段可能看到新值收到 update 事件包含前后值删除无法发现收到 delete 事件新增列需要人工改查询语句可能收到 DDL 事件实时性与轮询间隔一致毫秒级到秒级建议你在自己的测试环境中跑一遍把结果记录在文档里。这样团队内部做技术决策时就有客观依据而不是靠“我感觉”。7. 接口 API 与批量任务7.1 游标同步的批量任务设计游标同步天然适合批量任务。你可以把同步逻辑放到定时调度平台如 Airflow、Jenkins、xxl-job每次任务启动时读取持久化的游标处理完一批后更新游标。如果任务失败只需要从上次游标处重跑不会重复处理已提交的批次。批量任务最关键的是要把“处理逻辑”和“游标更新”放到同一个事务中。下面的伪代码展示了这个设计def process_batch(conn, target_conn, rows): # 在同一事务里写入目标库 target_conn.begin() try: for row in rows: insert_into_target(row) update_cursor(cursor) target_conn.commit() except Exception: target_conn.rollback() raise如果处理逻辑涉及外部 API比如调用第三方搜索引擎接口建议引入重试机制。游标更新可以等外部调用成功后再推进否则失败批次里的数据会永久丢失。7.2 CDC 的接口与消费者组CDC 连接器输出的事件通常不是 HTTP API而是消息队列中的主题topic。你可以让下游服务作为消费者消费这些事件。如果团队习惯 API 调用也可以在 CDC 事件消费端封装一层 HTTP 接口或者用 websocket 推送。下面是一个通用的 Python 消费者模板用于消费 Kafka 中的 CDC 事件from kafka import KafkaConsumer consumer KafkaConsumer( mysql-server.source_db.user_order, bootstrap_servers127.0.0.1:9092, auto_offset_resetearliest, enable_auto_commitFalse, group_iduser_order_syncer, value_deserializerlambda m: json.loads(m.decode(utf-8)), ) for message in consumer: event message.value op event[op] # c 表示 create u 表示 update d 表示 delete payload event[after] if op in (c, u) else event[before] print(f收到事件: {op} - {payload}) # 处理成功后手动提交 offset consumer.commit()注意这里需要根据你使用的 CDC 工具调整 topic 命名和消息结构。Debezium 默认的 topic 名是database.server.name database table而 Flink CDC 可能直接以 DataStream 方式输出不经过 Kafka。7.3 批量任务的失败重试建议无论是游标同步还是 CDC 消费批量任务都需要关注“重复消费”和“消息丢失”问题。游标同步重复如果一批数据写入目标库成功但游标更新失败下次会重复处理。因此目标表最好有唯一键或者更新操作使用ON DUPLICATE KEY UPDATE。CDC 重复因为 Kafka 至少一次语义下游处理必须支持幂等。更新写入目标库时可以使用数据库判断操作版本或者直接覆盖更新。失败重试给每条事件增加一个全局递增 ID如 binlog 位置目标库记录处理位置启动时从上次位置恢复。8. 资源占用与性能观察8.1 游标同步的数据库压力你可以用SHOW PROCESSLIST或performance_schema观察同步任务产生的查询。关键指标有两个扫描行数和锁等待。如果每次增量查询都扫描大量历史行说明索引设计不合理。比如WHERE id ?如果没有主键索引或者排序字段不匹配性能会直线下降。为了降低压力请确保游标字段上有索引主键天然满足查询使用覆盖索引减少回表批量大小要根据源库 IO 能力调整一般建议 500~2000 行轮询间隔不宜过低结合业务可接受延迟来判断8.2 CDC 的日志读取与内存占用CDC 连接器会持续读取 binlog 并解析事件。在 MySQL 中开启 binlog 本身会增加写入 IO如果 MySQL 实例同时处理大量写请求磁盘压力会明显上升。Debezium 或 Flink CDC 进程的内存占用与事件大小、批量发送量有关通常需要分配 2~4GB 堆内存给连接器具体以实际运行情况为准。观察日志解析趋势时可以关注binlog文件数量增长速度连接器的 lag 指标Kafka 消费组的current-offset与log-end-offset差值如果发现消费积压优先排查下游处理瓶颈而不是盲目增加连接器内存。8.3 性能调优要点游标同步调整批量大小、轮询间隔、选择合适游标字段必要时引入多线程并发读取多个分片。CDC开启并行消费增加连接器任务数但要注意保持分区内事件有序。如果源库是分库分表需要按分片键设计 topic 分区。网络传输如果跨机房同步需要关注带宽和延迟。CDC 传输的是紧凑格式通常比游标同步的整行查询更节省流量。9. 常见问题与排查方法问题现象可能原因排查方式解决方案游标同步漏数据游标字段不唯一或更新了游标字段检查游标字段是否唯一、是否可更新改用主键或结合 updated_at游标同步重复数据批次处理成功后游标未更新检查任务日志和游标存储将处理与游标更新放入同一事务游标同步无法感知删除设计缺陷确认业务需求是否需要删除事件增加软删除标记或迁移到 CDCCDC 连接器连接失败数据库用户权限不足检查用户是否能读取 binlog授予 REPLICATION SLAVE, SELECT 权限CDC 消费积压下游处理太慢查看消费者 lag 指标增加下游并行度优化目标库写入binlog 文件不存在日志过期被清理查看expire_logs_days和 LSN 对比调整日志保留策略或从全量快照重建更新事件无 before 字段配置了 binlog_row_imageMINIMAL查看 mysql 变量 binlog_row_image改为 FULL时区不一致导致时间游标错乱数据库和任务时区不同统一使用 UTC 时间并检查连接参数显式设置 serverTimezoneDDL 变更导致 CDC 任务停止连接器无法解析新结构查看连接器日志和 schema history升级连接器版本或手动处理 Schema 变更10. 最佳实践与使用建议从我的经验看选型前先回答三个问题业务是否需要捕获删除事件延迟要求是多少运维能力能不能支撑 CDC 组件这三个问题的答案基本能确定方案方向。如果你最终选择游标同步请先把以下事情做清楚设计一个不会进步的游标。优先使用主键如果需要处理更新引入updated_at加一个微调窗口防止时间相同导致漏数据。对目标库做幂等写入。一旦游标更新失败重复执行不会产生脏数据。明确删除的处理策略。如果业务有删除要么用逻辑删除标记要么在目标库定期执行全量对账否则数据会不一致。监控游标推进速率。出现长时间不推进、扫描行数暴涨立即报警。如果你选择 CDC则要注意提前设计消费幂等。CDC 事件至少一次语义目标库按主键 upsert 是基本要求。保存 binlog 位置。方便故障恢复时重新定位也方便做数据对比。测试 DDL 流程。上线前专门演练一次加列、减列、改类型的场景避免结构变更导致同步链路中断。如果只是需要最终一致不要引入 CDC。因为多一个组件就多一份故障率简单游标同步可能更合适。合规方面如果你的数据来自生产环境涉及用户隐私请确保所有同步链路都经过授权并且在传输和落地时做好加密、脱敏和审计。不要因为追求技术指标而忽略数据安全边界。11. 总结与下一步这篇内容最值得你记住的不是“CDC 比游标同步好”而是“很多数据增量问题其实不是单靠读取当前数据快照就能解决的”。游标同步适合简单的批量增量任务CDC 适合需要完整变更流和低延迟的系统。你需要先判断自己真的需要什么再决定要 build 到什么程度。如果你正在选型我建议先从一次小型的 POC 开始用一个测试库跑一遍第 6 节的功能测试记录结果。测试完你自然就知道两条路的差别了。最容易踩的坑不是技术细节而是对删除和更新历史事件的忽视。很多系统上线后才发现数据对不上就是因为没有提前设计好“被删除的旧数据如何处理”。先把这一点想清楚再动手写代码可以少走不少弯路。后续你可以继续探索的方向包括基于 CDC 的实时数仓搭建、游标同步与消息队列结合的分批方案、以及利用 ClickHouse 或 Iceberg 做数据湖增量同步。每一步都会对现有系统有实际的加固作用。