Maxwell:轻量 MySQL CDC、Bootstrap 与 Producer 恢复实战
一个推荐流服务曾出现过很难解释的空洞:MySQL 的订单状态从 paid 变成 shipped,Kafka topic 中却找不到对应更新。Maxwell 进程没有退出,复制延迟也很低,只有日志里一条很快被滚动覆盖的 producer timeout。重启服务没有补回消息,因为它已经继续处理后续 binlog;团队又用一次全表 bootstrap 修复,结果旧快照行覆盖了更新后的状态,空洞变成了回退。
事故的根因不是“Kafka 偶尔抖动”这么简单。Maxwell 同时维护 MySQL binlog 位点、表结构历史和 producer 异步确认,默认配置还允许某些 producer 错误只记日志后继续。bootstrap 又是一条插入历史行的并行数据流,消费者如果没有版本比较,就会把晚到的旧行当成最新事实。要安全使用这个轻量工具,必须把每一步状态、失败后是否前进、以及下游如何收敛说清楚。
先判断 Maxwell 是不是合适的那把工具
Maxwell 是一个读取 MySQL binlog 并把行变更编码成 JSON 的 daemon,内置 Kafka、Kinesis、Pub/Sub、RabbitMQ、Redis 等 producer。官方 GitHub 仓库 采用 Apache License 2.0,当前没有归档标记,并持续发布 v1.45.0;官网 Quick Start 的下载示例和 API 页面仍停在更早标签,因此安装时以仓库 Releases 为准、锁定具体版本,并把“网站内容滞后于发布包”纳入升级复核,而不是据旧页面误判项目停止维护。
它的优势是简单:不需要 Kafka Connect 集群,不引入统一 connector control plane,一条配置即可把 MySQL DML 变成易消费的 JSON;还内置按表 bootstrap 和 MySQL 中的 schema/position store。它的边界也同样清楚:核心围绕 MySQL 单源和单进程复制循环设计,不提供通用多数据库连接器平台、端到端校验、目标端事务、自动切流或成熟的跨机房 HA。
适合 Maxwell 的现场通常具备这些特征:数据源是 MySQL;输出主要是 JSON 事件流;任务数量不多;团队愿意自己治理 topic、消费者幂等和 schema 契约;对 bootstrap 的按表扫描与顺序语义有清晰理解。若需要数百连接器集中编排、跨多种数据库、schema registry、全量分片编排、savepoint 或复杂流式转换,Debezium/Kafka Connect、Flink CDC 或托管迁移服务通常更符合平台化目标。
这里最需要警惕的是 positions 与真实交付之间的距离。Maxwell 能根据 producer 回调管理处理进度,但如果配置选择“忽略 producer 错误”,日志中的失败不会自动变成可恢复消息。业务上必须把 failed counter、死信、目标水位和源位点联系起来,不能只看复制线程是否继续前进。
配好 MySQL 的日志、复制身份和状态库权限
官方 Quick Start 要求 binlog_format=ROW、启用 binlog、配置唯一 server_id,并给 Maxwell 复制与读取权限。官网 Compatibility 还说明运行时需要 JRE 11 或更高版本,并支持 MySQL 5.1、5.5、5.6、5.7、8;具体小版本、云数据库兼容性和认证插件仍要在目标环境验证。
# my.cnf
[mysqld]
server_id=201
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL
gtid_mode=ON
enforce_gtid_consistency=ON
log_replica_updates=ON
binlog_expire_logs_seconds=604800binlog_row_image=FULL 对 CDC 更友好。官方明确提示,MINIMAL 下 INSERT 可能没有默认值,UPDATE 的 data 可能只含主键和变化列,DELETE 通常只含主键。若下游把 Maxwell 的 data 当完整行,切到 MINIMAL 后会静默清空未出现字段。降低 binlog 体积必须与消费者 patch 语义一起评审。
Maxwell 使用 MySQL 做三件事:保存 schema 与 position、读取复制日志、捕获当前 schema。最简单部署让三个角色都指向 localhost。状态库账号需要管理 maxwell.*,复制账号需要读取业务表和 binlog。为了便于入门可先用同一账号,生产上应评估拆分 host、replication_host、schema_host 后的凭证边界。
CREATE USER IF NOT EXISTS 'maxwell_cdc'@'localhost'
IDENTIFIED BY '<MAXWELL_DB_PASSWORD>';
GRANT ALL ON maxwell.* TO 'maxwell_cdc'@'localhost';
GRANT SELECT, REPLICATION CLIENT, REPLICATION SLAVE ON *.*
TO 'maxwell_cdc'@'localhost';
CREATE DATABASE IF NOT EXISTS lab_orders
CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci;
CREATE TABLE lab_orders.orders (
id BIGINT PRIMARY KEY,
status VARCHAR(24) NOT NULL,
amount_cents BIGINT NOT NULL,
version BIGINT NOT NULL,
updated_at TIMESTAMP(6) NOT NULL
);
INSERT INTO lab_orders.orders VALUES
(5001, 'paid', 8800, 1, CURRENT_TIMESTAMP(6));
SHOW VARIABLES WHERE Variable_name IN
('log_bin', 'binlog_format', 'binlog_row_image', 'server_id',
'gtid_mode', 'enforce_gtid_consistency', 'log_replica_updates');
SHOW GRANTS FOR 'maxwell_cdc'@'localhost';ALL ON maxwell.* 是为了让 Maxwell 初始化和维护自己的元数据表,不应扩大到业务库。业务库通常只需 SELECT,复制权限则作用于全局。MySQL 8.4 仍以 REPLICATION SLAVE 允许复制客户端请求源端更新,以 REPLICATION CLIENT 允许读取 binary log 状态;不要用名字相近的 REPLICATION_SLAVE_ADMIN 代替数据面权限。如果组织要求更窄的表可见性,应配合 Maxwell filter 排除不可读表;官方 Deployment 提醒,权限看不到某些 schema 时可能出现 Can't find table,ignore_missing_schema 只能与严格过滤一起使用,不能作为忽略授权漂移的全局开关。
密码不应进入 shell history、镜像层或 Git。示例用占位符表示秘密;实际部署通过受控环境变量或只读配置文件注入,并开启 MySQL TLS。ssl=VERIFY_IDENTITY 比仅加密不校验主机身份更能抵御错误端点与中间人,但证书链和主机名必须先验证。
安装并以 stdout 建立最小可观察闭环
到 Maxwell Releases 下载锁定包。官网发布与主站文档可能不同步,因此先确认 tarball、变更日志和 JRE 基线,再写入组织制品库。下面使用当前稳定标签,不使用 latest。
MAXWELL_VERSION=1.45.0
curl -fL -o maxwell.tar.gz \
"https://github.com/zendesk/maxwell/releases/download/v${MAXWELL_VERSION}/maxwell-${MAXWELL_VERSION}.tar.gz"
mkdir -p ./maxwell-home
tar -xzf maxwell.tar.gz -C ./maxwell-home --strip-components=1
java -version
./maxwell-home/bin/maxwell --help | head -n 30第一次不要急着接 Kafka。stdout producer 能把 MySQL 权限、schema 初始化、binlog 解码和 JSON 输出问题与 broker 问题分开。把密码放在临时环境变量,再通过配置映射注入;调试结束后清除会话变量。
# config.properties
host=localhost
port=3306
user=maxwell_cdc
password=<MAXWELL_DB_PASSWORD>
schema_database=maxwell
client_id=orders_stdout
replica_server_id=9201
producer=stdout
filter=include: lab_orders.orders
gtid_mode=true
output_binlog_position=true
output_gtid_position=true
output_commit_info=true
output_primary_keys=true
output_primary_key_columns=true
output_schema_id=true
output_ddl=true
bootstrapper=asynccd ./maxwell-home
bin/maxwell --config=../config.properties这里同时设置 gtid_mode=true 与 output_gtid_position=true:前者让复制器按 GTID 定位,后者只决定 JSON 是否输出可用的 GTID position,并不会替 MySQL 开启 GTID。可启动前提是源端 gtid_mode=ON、enforce_gtid_consistency=ON,复制链路保留 GTID 历史并记录副本上的更新;切换场景中的候选主还必须包含 Maxwell 已消费事务的 GTID 集合。若实验环境不具备这些前提,应同时设为 false 并只验证 file/position,不能只删 gtid_mode 却保留“已验证 GTID 输出”的叙述。
启动成功应先看到元数据初始化、schema capture 和 binlog 连接,再在源表变更后出现 JSON;启用 GTID 的实验还应在事件中看到非空 GTID position,并用 SHOW VARIABLES 的结果证明源端前提。若停在连接重试,核对 host 匹配、认证插件、TLS 和端口;若报 server id 冲突,修改 replica_server_id;若报缺表结构,检查业务表 SELECT 权限与 filter。进程启动而没有事件,还要确认更新发生在启动读取位点之后。
Docker 也可运行官方镜像,但生产前需锁定不可变摘要、扫描依赖并持久化配置。状态本身保存在 MySQL maxwell 库,容器仍需稳定的 client id、replica server id 和网络身份。多个副本直接同时启动并不会自然变成安全 HA,后文会解释原因。
读懂 schema store、client id 与恢复位点
Maxwell 解析 row event 时必须知道当时表的列顺序、类型和主键。它把捕获到的 schema 及其演进保存在 schema_database 指定的库中,同时用同一库的 positions 表保存每个 client_id 的复制进度。配置项名称虽然是 schema_database,语义却是“Maxwell 状态库”;官方没有另一个 position_schema 配置项。把不存在的 position_schema=maxwell_position 写进 properties 不会把位点拆库,只会制造一条看似合理、实际未生效的配置。schema_id 是 Maxwell 自己跟踪的 schema 版本,不是 MySQL 原生 id。
这意味着 client_id 是恢复身份,不是随便写的实例名。两个不同任务使用同一 client id,会争用同一位置;同一任务重启时换 client id,则会像一条新订阅,起点和历史状态不再连续。官方要求同一主库上的每个 Maxwell 还使用唯一 replica_server_id,并且不能与任何 MySQL server id 冲突。
可以只读检查状态库,不要手工修改:
USE maxwell;
SHOW TABLES;
SELECT * FROM positions WHERE client_id = 'orders_stdout';
SELECT * FROM schemas ORDER BY id DESC LIMIT 5;
SELECT * FROM bootstrap ORDER BY id DESC LIMIT 5;
SELECT TABLE_NAME, COLUMN_NAME
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = 'maxwell'
AND TABLE_NAME IN ('positions', 'schemas')
ORDER BY TABLE_NAME, ORDINAL_POSITION;预期能在 positions 看到 client_id、server id、file/position 或 GTID 集合,并在 schemas 看到与 binlog position 绑定的结构快照/增量元数据。这个查询也能直接戳破 position_schema 误配:位点仍在 schema_database 指向的库里。状态库与业务 binlog 是一个恢复协议。若恢复了旧版 maxwell 库却没有足够旧 binlog,保存位置会指向不存在文件;若只保留 binlog却丢失 schema 历史,Maxwell 未必能用当前表结构正确解释旧事件;若克隆状态库后两边使用相同 client id,同时连接同一源,则会形成双生产者与位点竞争。
--init_position 能绕过 positions 指定 file/position,但官方明确把它标为危险,而且只能指向 Maxwell 已访问过的位置;它不支持写入 config.properties。把它当日常回退按钮会绕过 schema 对齐。真正需要回放时应复制状态库到隔离环境、确认目标位点对应 schema、输出到隔离 topic,并计算重复范围。
GTID 模式通过 gtid_mode=true 帮助主库切换后寻找事务身份。它依赖 MySQL 拓扑有一致 GTID 历史,并仍需要把 Maxwell 指向新主或稳定 VIP。GTID 解决“事务在哪条日志上”,不解决 producer 已确认与目标已应用之间的重复。
JSON 输出协议要被当成版本化契约
Maxwell 的 Data Format 中,DML 事件包含 database、table、type、ts、data,可选位置、server id、thread id、主键等。UPDATE 的 data 是更新后的行,old 只列变化字段的旧值;DELETE 的 data 是删除前行。BLOB、BINARY 和二进制字符串会 base64 编码,字符串统一输出 UTF-8。
一条启用位置与 schema id 的更新可能呈现为:
{
"database": "lab_orders",
"table": "orders",
"type": "update",
"ts": 1900000000,
"xid": 42,
"commit": true,
"position": "mysql-bin.000123:4567",
"schema_id": 18,
"primary_key": [5001],
"primary_key_columns": ["id"],
"data": {
"id": 5001,
"status": "shipped",
"amount_cents": 8800,
"version": 2,
"updated_at": "<T+N>"
},
"old": {
"status": "paid",
"version": 1
}
}xid + commit 可帮助重组事务,但不要把 xid 当永久全局 id;同一 server 生命周期之外未必稳定。position 有助于审计和去重,却在主库切换后需要结合 server/GTID 解释。主键字段应开启,因为 Kafka message key 和下游幂等通常依赖它们。没有主键的表仍能输出,但更新、删除和 bootstrap 收敛都缺少稳定身份,应在接入前整改或定义可靠业务键。
消费者契约至少规定:字段缺失与 null 的区别;decimal、unsigned、时间和零日期的处理;base64 字段如何解码;未知字段是否忽略;新事件 type 如何隔离;主键变更如何表达;事务何时可见。将 output_naming_strategy 从原字段名切到驼峰属于破坏性协议变更,不能只改 producer 配置。
用正向实验验证 DML、DDL 与事务结束
在 stdout 模式运行时提交下面的合成变更:
START TRANSACTION;
UPDATE lab_orders.orders
SET status='shipped', version=2, updated_at=CURRENT_TIMESTAMP(6)
WHERE id=5001;
INSERT INTO lab_orders.orders
VALUES (5002, 'paid', 12900, 1, CURRENT_TIMESTAMP(6));
COMMIT;
ALTER TABLE lab_orders.orders
ADD COLUMN channel VARCHAR(20) NOT NULL DEFAULT 'web';
UPDATE lab_orders.orders
SET channel='partner', version=3, updated_at=CURRENT_TIMESTAMP(6)
WHERE id=5001;预期输出包含同一事务的 update 与 insert,并在最后一行标出 commit=true;启用 output_ddl=true 后会出现 table-alter,包含 sql、变更前 old schema 和变更后 def schema;随后 DML 的 schema_id 应对应更新后的结构。若只看到 DML 而没有 DDL,先查配置与 ddl_kafka_topic 路由,不要据此认为 schema 没变化。
三层验证分别回答三个问题:源端最终状态是否正确,事件契约是否完整,目标是否按版本收敛。
SELECT id, status, amount_cents, version, channel
FROM lab_orders.orders ORDER BY id;
-- 下游示意
SELECT entity_id, source_position, source_schema_id, apply_status
FROM cdc_apply_log
WHERE stream_name='orders'
ORDER BY applied_at DESC LIMIT 20;目标端不能只计消息数。应检查 5001 最终 version=3、channel=partner,同一 source position 只有一次成功应用,未知 schema id 没有被静默跳过。若 DDL 已到达而消费者仍按旧字段白名单反序列化,应该在兼容层接受新字段或暂停,而不是丢弃整条事件。
Bootstrap 是一段受标记的历史流,不是普通快照
Maxwell 的 Bootstrapping 会对指定表执行 SELECT *,并把结果写入同一输出流。序列以 bootstrap-start 开始,随后每行一个 bootstrap-insert,再输出扫描期间积累的普通 INSERT/UPDATE/DELETE,最后发送 bootstrap-complete。
先创建独立配置并在写入秘密之前把权限收紧,再由秘密管理系统渲染内容。下面的文件只供 bootstrap 工具读取,不提交到 Git,也不与普通用户共享:
cd ./maxwell-home
umask 077
install -m 600 /dev/null ../bootstrap.properties# bootstrap.properties,文件权限必须为 0600
host=localhost
port=3306
user=maxwell_cdc
password=<MAXWELL_DB_PASSWORD>
schema_database=maxwell
client_id=orders_stdoutcd ./maxwell-home
test "$(stat -c '%a' ../bootstrap.properties)" = 600
bin/maxwell-bootstrap \
--config=../bootstrap.properties \
--database lab_orders \
--table orders \
--client_id orders_stdout \
--comment 'rebuild-search-projection'Maxwell bootstrap 工具支持 --config,命令行参数优先于配置文件。密码放进 0600 文件后不会出现在 shell history 或普通进程参数列表;生产中还应让文件只在作业期间存在,并由秘密系统轮换和销毁。client_id 要指向负责执行请求的 Maxwell 实例。--where 可限制行,但它直接进入查询条件,必须由受控作业模板生成,不能把外部输入拼接进去。bootstrap 表中的请求、开始和完成状态应纳入审计;未完成任务不能靠再次提交相同请求来假定幂等。
默认 bootstrapper=async 使用独立线程扫描。主复制线程继续处理其他表,但被 bootstrap 表在扫描期间发生的 binlog 事件会在进程内排队,等历史行发完后再发送,最后才 complete。官方配置样例直接提示 async 可能因缓冲事务造成数据丢失;它不是一条可假设为持久化的磁盘队列。sync 使用复制线程扫描,期间所有 binlog 事件阻塞,逻辑更直观但会让全局复制 lag 上升。大表、长事务和高写入表必须先压测队列峰值、JVM 堆、源端 lag 与崩溃恢复,而不能用扩大堆内存掩盖 producer 长期跟不上。
消费者应把 bootstrap-insert 视为基线候选,并以业务版本、更新时间或源事件序与后续增量收敛。不能因为 bootstrap 消息晚到就覆盖当前行。bootstrap-complete 只证明 Maxwell 完成扫描和发送流程,不证明 Kafka 所有副本、消费者或目标已完成,也不证明行数与源库一致。
用反向实验暴露 Bootstrap 覆盖和 Producer 空洞
第一组反向实验验证 bootstrap 与并发更新。启动 async bootstrap 后,在扫描未结束时把 5001 从 version 3 更新到 version 4。正确流序应先有 bootstrap-start 和历史行,再有排队的普通 update,最后 complete;目标端最终必须保留 version 4。如果目标按“最后收到的 bootstrap-insert”无条件 upsert,就可能回到 version 3。
语义证据可以记录为:
bootstrap-start table=lab_orders.orders
bootstrap-insert id=5001 version=3
update id=5001 old.version=3 data.version=4
bootstrap-complete table=lab_orders.orders
target id=5001 version=4 reconciliation=matched再做一次崩溃反例:看到若干 bootstrap-insert 后终止隔离环境中的 Maxwell,再用同一 client_id 重启。官方定义的恢复行为是整段 bootstrap 重新执行,而不是从 inserted_rows 断点续扫,因此预期会再次出现 bootstrap-start 和已经发送过的行。消费者必须把重复历史行收敛到同一业务版本;同时核对崩溃窗口内那条 version=4 普通更新最终仍出现。若未出现,说明 async 缓冲窗口已经形成空洞,应停止接管并从源表与 binlog/审计记录补偿。
失败任务不会因为再次提交同一 bootstrap 请求就天然幂等。决定放弃时,先保留请求 id、已发主键范围和目标对账证据,再按官方故障处理语义将旧请求 is_complete=1 或删除旧行,随后才创建新的受审计请求。直接清空整个 bootstrap 表会同时破坏其他任务的恢复证据。
第二组反向实验验证 producer 失败。Maxwell Reference 中 ignore_producer_error 默认是 true:Kafka、Kinesis、Pub/Sub 的发布错误可能只记录日志后继续。对于迁移或关键投影,这个默认值通常过于宽松,应显式设为 false,让非特例发布错误终止进程并阻止问题被低噪声掩盖。
在隔离环境可临时撤销测试 producer 对 topic 的写权限,或把 broker 指向不可达端口;随后写入一条唯一版本更新。预期是 messages.failed 增加、健康检查失败或进程退出,源端与 Maxwell 位点差开始扩大。恢复权限并重启后,应验证该 source position 是否重新发布,并以目标幂等记录确认只生效一次。
producer=kafka
kafka.bootstrap.servers=localhost:9092
kafka_topic=cdc.orders
dead_letter_topic=cdc.orders.dead-letter
ignore_producer_error=false
producer_ack_timeout=30000
kafka.acks=all
kafka.enable.idempotence=true
kafka.retries=10producer_ack_timeout 是针对异步 producer 没有返回成功或失败的启发式超时,不是端到端交付保证。dead_letter_topic 在 Kafka producer 中遇到发布行错误时写入只含主键的 skeleton row,尤其 RecordTooLargeException 可走此路径;它不能重建被截掉的大字段,因此死信必须触发源端按主键重取或专门修复,而不能计作成功消息。
即使 Kafka producer 开启 idempotence,也只减少 producer 重试在 Kafka 内造成的重复,不会把 MySQL 事务、Kafka 消息和目标数据库事务合并成一个原子提交。下游依然需要主键、version/source position 和幂等事务。
Producer、分区与顺序决定下游能否收敛
官方 Producers 说明,Kafka 配置使用 kafka. 前缀透传给客户端。可靠性基线通常包括 acks=all、足够重试和 topic 的 min.insync.replicas。Maxwell 还能按 database、table、primary key、transaction id、column 或 random 选择分区键。
按主键分区可以保持同一订单的更新顺序并获得较好并行度;按 transaction id 能让同一事务聚合,却可能使同一实体跨事务落到不同分区;按 table 或 database 扩大顺序域但形成热点;random 吞吐均匀却几乎放弃实体顺序。选型要从消费者的不变量倒推,而不是只追求均匀。
Maxwell 启动时读取 topic 分区数。扩分区后 hash 到 partition 的结果可能变化,同一主键的新旧事件可处于不同分区,跨分区消费没有全局顺序。需要严格实体顺序时,应在扩容前暂停生产、排空旧分区、记录切换水位,或让消费者以 version 拒绝旧事件。
大消息是另一类边界。BLOB base64 会比原始二进制更大,完整行 UPDATE 也会重复发送未变化字段。提高 Kafka max.request.size 前要同时核对 broker 和 consumer 上限、网络与堆内存;更可控的做法是过滤大字段、输出对象存储引用,或为 LOB 表建立单独链路。死信 skeleton 只有主键,不能被当成完整替代事件。
重启、主库切换和 HA 不应被同一个开关概括
普通重启依靠 maxwell.positions 与 schema store 继续。测试时要在 producer 成功后和 position 持久化前后分别杀进程,观察是否有限重放;目标以 version 或 source position 去重。若所需 binlog 已过期,positions 再完整也无法恢复,必须从归档日志或新 bootstrap 建立基线。
主库切换优先使用 MySQL GTID,并启用 gtid_mode=true。官方 Deployment 说明 Maxwell 仍需被重新指向新主或通过稳定 VIP 连接。非 GTID 的 master_recovery 通过写入和回看 heartbeat 在新主日志中寻找近似位置,但官方 Compatibility 将其标为 alpha,提示高活跃服务器可能重复最多约一秒数据,而且它与分离的 schema/replication host 不兼容。
Maxwell 自带基于 jgroups-raft 的 client-side HA,但官方 High Availability 同样明确标为 experimental、alpha quality,需要至少三个节点。它提供 leader election,不代表 position、schema store、broker 和 MySQL 都具备同等级容灾,更不能替代 split-brain 演练。
<!-- raft.xml 中的核心成员示意;网络地址需按受控环境配置 -->
<raft.RAFT members="A,B,C" raft_id="${raft_id:undefined}"/># 三个隔离节点分别使用唯一 member id,共享一致配置和状态库
bin/maxwell --config=config.properties --ha --raft_member_id=A
bin/maxwell --config=config.properties --ha --raft_member_id=B
bin/maxwell --config=config.properties --ha --raft_member_id=C生产决策不能因为“有 HA 参数”就采用它。若团队没有能力测试网络分区、双 leader、状态库不可用和成员滚动升级,更稳妥的方案往往是单 active + 外部编排冷/温备:同一时刻只允许一个实例持有复制身份,故障后由带 fencing 的控制器启动备用。RTO 会稍长,但行为更容易证明。
HA 演练应保存故障前 source position、leader 身份、producer 成功水位、故障后起始位置、重复数和目标最终版本。只证明新进程起来了,没有证明数据连续。
DDL 与 Schema 漂移既是能力也是风险
output_ddl=true 后,Maxwell 可输出 database-create、database-alter、database-drop、table-create、table-alter、table-drop 等事件。table-alter 包含旧结构 old、新结构 def、SQL 和位置。它内部也用这些 DDL 更新 schema store,以便后续 row event 按新列布局解码。
DDL 输出仍不是 schema registry。消费者必须决定:新增可空列是否可忽略,重命名是否视作删加,decimal 精度收缩是否阻断,主键变化如何迁移 Kafka key,字符集变化是否需要重编码。未知 DDL 应进入隔离状态并暂停相关表,不要只记录 warning 后继续把错列数据写入目标。
源端变更最好采用 expand-contract:先让消费者接受新旧 JSON;再执行向后兼容 DDL;观察新 schema id;迁移数据并切换生产者;最后删除旧字段。ddl_kafka_topic 可把 DDL 放进独立 topic,但这会制造两个 topic 间的顺序协调问题。消费者必须以 source position/schema id 建立屏障,不能假设先消费控制 topic 再消费数据 topic 就天然有序。
recapture_schema 能重新捕获最新 schema,却不是修复任意历史错位的按钮;它不在 properties 中提供,且会改变后续解释基线。使用前应复制状态库、冻结相关表、定位第一个错误位点,并评估旧事件是否需要隔离回放。
项目接入要把目标端幂等做成第一等状态
一个可靠消费者至少保存 stream_name + entity_key + source_position/schema_id + business_version + apply_status。处理一条更新时,先比较 business version;新版本才更新投影,并在同一目标事务写入 apply log。重复 source position 返回成功但不重复副作用,旧 version 记录为 stale,未知 schema 进入 quarantine。
CREATE TABLE cdc_apply_log (
stream_name VARCHAR(64) NOT NULL,
event_key VARCHAR(160) NOT NULL,
source_position VARCHAR(160) NOT NULL,
schema_id BIGINT NULL,
business_version BIGINT NULL,
apply_status VARCHAR(24) NOT NULL,
PRIMARY KEY (stream_name, event_key, source_position)
);
-- 业务投影更新应带版本保护
UPDATE order_projection
SET status = :status, version = :incoming_version
WHERE id = :id AND version < :incoming_version;bootstrap-insert 和普通 insert/update 走同一收敛函数,但保留不同 event type 便于审计。DELETE 需要 tombstone 策略:物理删除、软删除或版本化墓碑必须统一,否则晚到 bootstrap 会把已删除行复活。副作用型消费者,例如发短信或扣库存,不能仅靠数据库 upsert;应使用 inbox/outbox 或独立事件身份保证只触发一次。
项目接入门槛应包括契约测试:缺少可选字段、出现新字段、MINIMAL row image、不合法零日期、base64 BLOB、主键变化、DDL、新 bootstrap type、重复与乱序。只用一条正常 insert 测试,无法证明系统能承受恢复现场。
权限、敏感数据与加密不能只盯数据库密码
Maxwell JSON 默认包含完整行,新旧值还可能同时出现;错误日志、stdout、死信、bootstrap comment 和监控诊断都可能泄漏敏感信息。stdout 只应用于合成数据和隔离环境,不能把生产行打印进集中日志。开启 output_row_query 还会复制原始 SQL,应默认关闭。
Maxwell 提供 encrypt=data 或 encrypt=all,但加密钥匙若与配置放在同一文件,风险并未降低。加密还会影响 Kafka key、过滤、调试和下游路由。更清晰的方案通常是 TLS 保护传输、broker 存储加密、严格 ACL,以及在受控转换层按字段令牌化;只有明确威胁模型要求消息级加密时,再引入独立 KMS、密钥版本和轮换协议。
HTTP metrics 默认绑定所有地址,启用时应设置 http_bind_address=localhost 或受控管理网,并通过反向代理鉴权。http_config=true 会开放动态修改 filter 的 /config,错误或越权 PATCH 可能瞬间扩大敏感表订阅或造成漏数;没有强认证、审计和配置基线时不要开启。
Kafka 凭证、云 producer 凭证和数据库密码分别最小授权。死信 topic 因为包含主键且指向失败数据,访问级别不应低于主 topic。bootstrap 作业权限应单独控制,因为一次全表 SELECT 会绕过日常增量速率限制并大规模复制数据。
容量、监控和故障证据要围绕最慢环节设计
Maxwell 是单线程 replicator:binlog 事件由一个线程捕获并逐条交给 producer;异步 producer 可并发在途,但持续吞吐仍受解析、序列化、producer ack 和状态持久化共同限制。官方 bootstrap 文档也明确说明 sync 模式会阻塞这一主线程。它适合轻量链路,不应在没有压测的情况下承接无限表和超大事务。
容量估算至少包括峰值行事件率、完整行 JSON 膨胀、base64 膨胀、producer 在途数、bootstrap 扫描速率、bootstrap 期间被修改表的排队量、Kafka 副本与保留、MySQL binlog 保留和 schema store 增长。buffer_memory_usage 默认使用 JVM 最大堆的一部分,扩大它只能延后背压,不能修复 producer 长期低于输入速率。
官方 Monitoring 暴露 messages.succeeded、messages.failed、row.count、replication.lag、inflightmessages.count、publish time/age 等指标,并提供 /metrics、/prometheus、/healthcheck、/ping。/ping 只证明 HTTP 线程活着;/healthcheck 对近期消息失败更敏感;业务接管还需比较源位点、Maxwell position、Kafka end offset、消费者 offset 和目标应用水位。
metrics_type=http
http_bind_address=localhost
http_port=8080
http_diagnostic=false
metrics_jvm=true
metrics_age_slo=60curl -fsS http://localhost:8080/healthcheck
curl -fsS http://localhost:8080/prometheus | grep -E \
'messages_(succeeded|failed)|replication_lag|inflightmessages'告警要能指导行动:failed counter 非零立即阻断接管;inflight 持续上升说明 producer ack 跟不上;replication lag 上升但 Kafka producer 正常,可能是大事务或 bootstrap 阻塞;lag 低但目标水位落后,故障在下游。监控保留 source position 和事件计数,不采集完整行。
清理、回滚、升级与团队治理决定能否长期使用
每个 Maxwell 任务应登记 owner、client id、replica server id、源表/filter、schema database、producer/topic、分区策略、输出字段开关、bootstrap 模式、binlog 保留、敏感级别、延迟预算和退出条件。client id 与 replica id 要由中央登记分配,避免团队各自复制默认 maxwell 与 6379。
升级时先用独立 client id 和隔离 topic 做影子读取,比较 DML 类型、position、schema id、DDL def、主键 key 和特殊类型编码。影子任务从新起点开始会缺少历史,不能只比较总条数;应在共同水位之后按源位点和主键对账。不要让新旧版本同时写一个非幂等 topic 再期待消费者自动判断来源。
回滚程序前保留旧包、旧配置、schema store 备份、最终 position 和 topic 契约。若新版本已经升级内部元数据 schema,旧程序是否可读必须通过副本演练;直接把生产状态库降级可能造成更深损坏。数据错误则按 source position 和目标 apply log 修复,不能用“回滚二进制”代替数据补偿。
退出链路时,先停止或记录源端最终写水位,等待 Maxwell、Kafka 和目标全部追平,完成分块 checksum 与业务不变量校验;关闭 producer;保存必要审计;撤销 MySQL 复制账号、Kafka ACL 和云凭证;按回切窗口保留 schema/positions 后再清理 maxwell 元数据、bootstrap 请求、topic 与死信。删除状态库不是第一步。
最终的选型结论应克制:Maxwell 是维护活跃、部署轻、JSON 直出的 MySQL CDC daemon,bootstrap 与 schema store 让小规模链路很高效;默认 producer 错误策略、单线程复制、alpha 级 HA 和缺少端到端迁移控制面,则要求团队补足幂等、校验、切流和治理。接受这些边界时它很好用,试图把它扩成万能迁移平台时,复杂度会悄悄转移到脚本和人工记忆里。
