Flink CDC Pipeline:从全量增量接续到 Schema 演进与幂等落库
作业恢复成功了,目标表为什么还是多出重复数据
一次 MySQL 到分析库的在线迁移已经追平,团队在发布窗口重启了 Flink 集群。控制台显示作业从 checkpoint 恢复成功,源端 lag 也重新下降,但目标表的订单金额突然翻倍。第二天又发生一次 DDL:源表新增可空列,作业这次直接全局失败。两个现象看似无关,根因却都来自同一个误解:Flink 能恢复算子状态,不等于任意 sink 都能把外部副作用恢复成恰好一次;Pipeline 能捕获 SchemaChangeEvent,也不等于任意目标系统都能原样执行每一种 DDL。
Flink CDC Pipeline 把数据库的快照与变更日志组织成 Flink 作业。source 保存快照 split、日志位点和表结构,SchemaOperator 协调结构事件,sink 把数据与结构变化提交到外部系统。checkpoint 是运行中的一致性状态切片,savepoint 是人为触发、用于受控升级或迁移的长期状态工件。目标端是否重复,还取决于 sink 是事务提交、幂等 upsert,还是无条件 append。
因此第一证据不能只是 JobManager 上的 RUNNING。要同时回答:作业从哪个 checkpoint/savepoint 恢复,source 位点是否连续,SchemaOperator 正在等待什么,sink 在故障前后的事务或批次如何标识,目标表用什么键收敛重放。缺一项,“恢复成功”都可能只是计算层成功。
一条 Pipeline 里真正流动的是三类状态
Flink CDC 的 YAML 由必需的 source、sink、pipeline 和可选的 route、transform 组成。source 不只输出行,还输出建表、加列、改类型等结构事件;route 决定源表与目标表的映射;transform 改变字段、过滤条件或表结构;pipeline 放全局并行度、时区、Schema 行为和稳定的 operator UID。
运行时可以把状态分成三类。第一类是 source 状态:哪些 snapshot chunk 已完成、binlog/WAL 读到哪里。第二类是 Flink 一致性状态:各算子在 checkpoint barrier 对齐时保存的状态和 channel 数据。第三类是外部提交状态:sink 已经预提交、提交或无法撤销的批次。前两类由 checkpoint 恢复,第三类必须由 sink 协议配合。
快照阶段把大表拆成 chunk 并行读取,日志读取器同时保住快照期间的增量,待各 chunk 完成后切入纯增量阶段。checkpoint 记录的是这组中间状态的共同切片,而不是“数据库在某个墙上时间的拷贝”。恢复时,已完成 chunk 不应全部重扫,日志读取从状态中的位点继续;但 checkpoint 周期之间已经送到外部系统的记录可能重放,sink 必须识别事务、批次或业务键。
先锁定 3.6 发布线和兼容矩阵
稳定部署应使用 Apache Flink CDC 3.6.x 发布文档、同一发行线的发行包与 connector JAR。其Pipeline connector 兼容矩阵列出 3.6.x 只支持 Flink 1.20.* 或 2.2.*,Pipeline source 包含 MySQL、PostgreSQL 与 Oracle,sink 包含 Doris、StarRocks、Paimon、Kafka、Elasticsearch、OceanBase、MaxCompute、Iceberg、Fluss 与 Hudi。矩阵外的组合不能因为类能加载就视为受支持。
发布文档域名中虽然包含 nightlies.apache.org,flink-cdc-docs-release-3.6 指向已发布的 3.6 文档线;3.7-SNAPSHOT 是未发布文档,不能作为生产参数基线。Flink CDC 3.6.0 只支持 JDK 11 或更高版本,项目主体采用 Apache License 2.0。MySQL Connector/J 的 GPLv2 许可与 Apache 项目分发不兼容,因此官方文档要求单独提供该驱动;企业制品清单既要记录 Flink CDC connector,也要记录 JDBC 驱动的来源与许可证。每次升级都要把 JDK、Flink、Flink CDC distribution、source connector、sink connector、JDBC 驱动和外部系统版本当作一个兼容单元。
本地验证可用 Flink standalone;共享测试环境可以跑 standalone、YARN 或 Kubernetes。生产形态的选择主要改变资源调度、升级和故障域,不改变 Pipeline 的 source/sink 契约。下面固定使用 Flink 1.20.* 与 Flink CDC 3.6.0,需要从官方发布页取得发行包和 3.6.0 connector JAR,并单独准备兼容的 MySQL Connector/J。下载后校验签名或 checksum;所有 JobManager/TaskManager 节点必须看到同一组不可变 JAR。
# 解压目录仅作示例,版本必须与兼容矩阵一致
tar -xzf flink-1.20.x-bin-scala_2.12.tgz
tar -xzf flink-cdc-3.6.0-bin.tar.gz
export FLINK_HOME="$PWD/flink-1.20.x"
export FLINK_CDC_HOME="$PWD/flink-cdc-3.6.0"
cp flink-cdc-pipeline-connector-mysql-3.6.0.jar "$FLINK_CDC_HOME/lib/"
cp flink-cdc-pipeline-connector-doris-3.6.0.jar "$FLINK_CDC_HOME/lib/"
cp flink-cdc-pipeline-connector-kafka-3.6.0.jar "$FLINK_CDC_HOME/lib/"
cp mysql-connector-java-8.0.27.jar "$FLINK_CDC_HOME/lib/"
ls -1 "$FLINK_CDC_HOME/lib"实际构件名、外部系统版本与传递依赖以官方 3.6 connector 下载页及各 connector 页面为准。企业代理环境应从受控制品库同步并保留来源、哈希和许可证清单,不让每个节点临时访问公网下载。JAR 不应通过共享可写卷被运行中替换;升级时创建新镜像或不可变制品,再用 savepoint 做受控迁移。
在第一次提交有状态作业前,先配置 checkpoint。否则作业即使已经写出数据,也没有可用于自动恢复的持久状态锚点。flink-conf.yaml 中的路径必须被所有 JobManager 和 TaskManager 访问;下面的 file:// 只适合所有进程共享同一本地文件系统的单机实验,分布式环境应换成受控对象存储或高可用文件系统:
state.backend: rocksdb
execution.checkpointing.interval: 30s
execution.checkpointing.mode: EXACTLY_ONCE
execution.checkpointing.timeout: 10min
execution.checkpointing.max-concurrent-checkpoints: 1
execution.checkpointing.num-retained: 3
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION
execution.checkpointing.storage: filesystem
execution.checkpointing.dir: file:///var/lib/flink/checkpoints
execution.checkpointing.savepoint-dir: file:///var/lib/flink/savepointscheckpoint 间隔越短,潜在重放窗口越小,但 barrier、状态上传和 sink 提交更频繁;间隔越长,开销下降,故障后重放和追平时间增加。超时不能盲目调大,应先判断是 source 背压、状态过大、存储延迟,还是 sink 预提交阻塞。修改配置后再启动 standalone 集群,并先确认 Flink 进程对两个目录有读写权限。
从 MySQL 到 Doris 写出第一条可审查 Pipeline
下面沿用官方 3.6 standalone 示例的 MySQL 到 Doris 组合,数据库和凭证均为合成值。MySQL 必须启用 ROW 格式 binlog,并将行镜像设为 FULL,使 UPDATE/DELETE 的行变化具备连接器预期的信息。先在 MySQL 配置文件中设置并重启实例,再核对实际值:
[mysqld]
server-id=5400
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULLSHOW GLOBAL VARIABLES WHERE Variable_name IN
('log_bin', 'binlog_format', 'binlog_row_image');
CREATE DATABASE IF NOT EXISTS sales_db;
CREATE TABLE sales_db.orders (
id BIGINT PRIMARY KEY,
customer_ref VARCHAR(64) NOT NULL,
amount DECIMAL(12,2) NOT NULL,
status VARCHAR(24) NOT NULL,
updated_at TIMESTAMP(3) NOT NULL
);
CREATE USER 'cdc_reader'@'10.42.16.%' IDENTIFIED BY '<MYSQL_CDC_PASSWORD>';
GRANT SELECT ON sales_db.* TO 'cdc_reader'@'10.42.16.%';
GRANT SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT
ON *.* TO 'cdc_reader'@'10.42.16.%';
SHOW GRANTS FOR 'cdc_reader'@'10.42.16.%';SELECT 支撑快照,SHOW DATABASES 支撑元数据发现,REPLICATION SLAVE 与 REPLICATION CLIENT 支撑 binlog 读取和位点查询;默认启用增量快照时不需要为了旧式全局锁额外授予 RELOAD。若捕获多个数据库,逐库授予 SELECT,不要为省事把所有业务库读权交给一个账号。
Doris 侧先由管理员创建目标库和受限账号。Pipeline 需要自动建表、写入和演进结构时,账号至少要在目标库具备 CREATE_PRIV、LOAD_PRIV 和 ALTER_PRIV;下面额外授予 SELECT_PRIV 用于目标对账,不授予全局管理权限:
CREATE DATABASE IF NOT EXISTS ods_sales;
CREATE USER 'cdc_writer'@'10.42.16.%' IDENTIFIED BY '<DORIS_CDC_PASSWORD>';
GRANT SELECT_PRIV, LOAD_PRIV, ALTER_PRIV, CREATE_PRIV
ON internal.ods_sales.* TO 'cdc_writer'@'10.42.16.%';
SHOW GRANTS FOR 'cdc_writer'@'10.42.16.%';从实际 JobManager/TaskManager 网段分别登录 MySQL 与 Doris,执行只读探针,才能同时验证账号、主机匹配、网络和权限,而不是只看管理员会话中的 SHOW GRANTS:
mysql -h mysql-source -u cdc_reader -p \
-e "SELECT COUNT(*) FROM sales_db.orders; SHOW MASTER STATUS;"
mysql -h doris-fe -P 9030 -u cdc_writer -p \
-e "SELECT CURRENT_USER(); SHOW GRANTS; USE ods_sales; SHOW TABLES;"把密码放入部署平台 Secret 或先渲染到权限受限的临时文件,不提交到 Git。
source:
type: mysql
name: sales-mysql
hostname: localhost
port: 3306
username: cdc_reader
password: <MYSQL_CDC_PASSWORD>
tables: sales_db.orders
server-id: 5401-5404
server-time-zone: UTC
scan.startup.mode: initial
scan.incremental.snapshot.chunk.size: 8096
heartbeat.interval: 30s
schema-change.enabled: true
sink:
type: doris
name: sales-doris
fenodes: localhost:8030
username: cdc_writer
password: <DORIS_CDC_PASSWORD>
table.create.properties.replication_num: 1
route:
- source-table: sales_db.orders
sink-table: ods_sales.ods_orders
pipeline:
name: sales-orders-cdc
parallelism: 2
local-time-zone: UTC
schema.change.behavior: exception
operator.uid.prefix: sales-orders-v1tables 是正则捕获表达式,点号是库表分隔符;批量捕获时要转义真正属于名称的点。server-id 范围应至少覆盖 source reader 的并发需求,并且不能与同一 MySQL 集群上的其他复制客户端冲突。scan.startup.mode=initial 先做快照再接 binlog;latest-offset 只读取启动后的变化,不能替代历史基线。chunk.size 增大可减少 split 数量,但会延长单块执行时间并增大失败重做与内存压力。
MySQL source 的 exactly-once 边界取决于分片键是否稳定。默认使用主键的第一列切 chunk;没有主键时必须显式设置非空 scan.incremental.snapshot.chunk.key-column。如果该非主键列在快照期间被更新,同一行可能跨 chunk 出现,官方只承诺至少一次,目标端必须按真正业务键收敛。scan.incremental.snapshot.backfill.skip=true 也会把快照窗口内变化留到日志阶段重放,官方明确提示只保证至少一次;不要为了缩短 backfill 把它当无损加速开关。heartbeat.interval 默认 30s,用于推进可恢复 binlog 位点,设为 0s 会禁用。
schema.change.behavior=exception 故意选择保守策略:遇到结构变化先失败,由变更流程判断目标端是否支持。等团队完成 Schema 演练后再选择 evolve 或 lenient,比默认吞下风险更容易建立证据。operator.uid.prefix 保持算子 UID 稳定,便于有状态升级和从 savepoint 映射旧状态;改 route、transform 或拓扑后,仍需验证状态是否可兼容,稳定 UID 不是万能迁移器。
提交到已经启动的 standalone Flink 集群:
"$FLINK_HOME/bin/start-cluster.sh"
"$FLINK_CDC_HOME/bin/flink-cdc.sh" mysql-to-doris.yaml \
--flink-home "$FLINK_HOME"
curl -fsS http://localhost:8081/jobs/overview3.6 standalone 部署文档的成功回显包含 Pipeline has been submitted to cluster 与 Job ID。回显只证明提交成功;还要在 Flink Web UI 或 REST 中确认 job 进入 RUNNING、checkpoint 能持续完成、source 从 snapshot 转为 stream reading,并在目标端核对行状态。
正向实验:让全量、增量和恢复形成证据链
在提交作业前准备源表与一条基线记录,提交后执行更新、新增和删除:
INSERT INTO sales_db.orders VALUES
(1001, 'C-LAB-01', 88.50, 'CREATED', CURRENT_TIMESTAMP(3));
UPDATE sales_db.orders
SET status='PAID', updated_at=CURRENT_TIMESTAMP(3)
WHERE id=1001;
INSERT INTO sales_db.orders VALUES
(1002, 'C-LAB-02', 19.90, 'CREATED', CURRENT_TIMESTAMP(3));
DELETE FROM sales_db.orders WHERE id=1002;预期目标 ods_sales.ods_orders 最终只有 id=1001,状态为 PAID;id=1002 被删除。验证时记录源库主键集合与业务聚合、Flink Job ID、最近成功 checkpoint ID、目标端主键集合和聚合结果。只比较 count 会漏掉“一行被错误更新、另一行刚好丢失”的抵消。
然后等一个 checkpoint 成功,记录 checkpoint 路径与 ID。自动故障恢复与取消后重提是两条不同路径,不能写成一个动作。
自动恢复实验先确认 standalone 至少有一个可重新启动或可接管 slot 的 TaskManager,再停止当前 TaskManager:
"$FLINK_HOME/bin/taskmanager.sh" stop
"$FLINK_HOME/bin/taskmanager.sh" start
curl -fsS http://localhost:8081/jobs/overview预期同一个 Job ID 经历 RESTARTING 后回到 RUNNING,Flink 自动选择最近成功 checkpoint,source 不重新做整表快照。若没有可用 TaskManager,作业会等待资源;这不等于状态丢失,启动 TaskManager 后应继续恢复。恢复后再更新 id=1001,检查新 checkpoint 能持续完成,目标最终状态与源端一致。
若主动取消作业,Flink 不会自动把已进入 CANCELED 的作业拉起。由于前面配置了 RETAIN_ON_CANCELLATION,可记录外部化 checkpoint 的 _metadata URI,然后显式用 -s 提交同一 Pipeline:
"$FLINK_HOME/bin/flink" cancel <JOB_ID>
"$FLINK_CDC_HOME/bin/flink-cdc.sh" \
-s file:///var/lib/flink/checkpoints/<JOB_ID>/chk-<N>/_metadata \
mysql-to-doris.yaml --flink-home "$FLINK_HOME"受控升级更适合先停止并生成 savepoint,再用同一个 -s 入口恢复:
"$FLINK_HOME/bin/flink" stop \
--savepointPath file:///var/lib/flink/savepoints <JOB_ID>
"$FLINK_CDC_HOME/bin/flink-cdc.sh" \
-s file:///var/lib/flink/savepoints/savepoint-<ID> \
mysql-to-doris.yaml --flink-home "$FLINK_HOME"-s 可以指向 savepoint 目录,也可以指向 checkpoint/savepoint 的 _metadata;路径必须仍可被所有执行节点读取。恢复窗口内可能重放状态快照之后、故障或取消之前的记录;目标是主键模型并按 key upsert 时,重复可以收敛。生产放行应在目标版本、目标 connector 和同规格网络上留下实际 checkpoint、TaskManager 恢复、取消后 -s 恢复与业务键对账证据。
反向实验:用官方 Kafka sink 重放同一状态
不需要编写一个无法复现的测试 sink。使用 Flink CDC 3.6 官方 Kafka Pipeline connector,把相同 source 的 sink 改到隔离 topic;Kafka topic 是 append 日志,不会因为业务主键相同自动覆盖旧记录:
sink:
type: kafka
name: replay-kafka
properties.bootstrap.servers: localhost:9092
topic: cdc-replay-lab
partition.strategy: hash-by-key
key.format: json
value.format: debezium-json先创建一个 cleanup.policy=delete 的实验 topic,提交作业并等待快照结束。随后生成一个中间 savepoint S0,让原作业继续运行;在 S0 之后插入唯一实验键,等消费者看见第一条记录,再取消原作业并从同一个 S0 恢复:
kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic cdc-replay-lab --partitions 1 --replication-factor 1 \
--config cleanup.policy=delete
"$FLINK_HOME/bin/flink" savepoint <JOB_ID> file:///var/lib/flink/savepoints
# 记下返回的 S0 URI,然后在 MySQL 执行:
# INSERT INTO sales_db.orders VALUES
# (9001, 'C-REPLAY-01', 9.01, 'CREATED', CURRENT_TIMESTAMP(3));
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
--topic cdc-replay-lab --from-beginning
"$FLINK_HOME/bin/flink" cancel <JOB_ID>
"$FLINK_CDC_HOME/bin/flink-cdc.sh" -s <S0_URI> \
mysql-to-kafka-replay.yaml --flink-home "$FLINK_HOME"原作业必须先取消,不能让两个实例从同一状态同时读取同一 source。恢复后再次从头消费 topic,按 key 筛出 id=9001;预期会看到两个不同 Kafka offset 的 create 记录:第一条是原作业已经追加的外部副作用,第二条是从 S0 恢复后重放的同一源变更。这个受控步骤证明“状态快照不会撤销已经写入的 append 目标”;随机故障是否重放以及重放多少,仍由最后成功状态、sink 提交协议和故障时刻共同决定。实验结束后停止作业并删除隔离 topic。
修复有三条常见路线。目标支持事务时,sink 将记录预提交到与 checkpoint 绑定的事务,在 checkpoint 完成通知后提交,恢复时中止悬挂事务。目标支持主键 upsert 时,使用稳定业务键,并用源版本或更新时间拒绝旧重放。只能 append 时,在落地层保存稳定事件 ID、source partition/offset 或业务版本,通过唯一约束和去重压实收敛。事件 ID 必须跨重启稳定;运行时随机 UUID 只会把同一重放变成两个不同事件。
不要仅依据“Flink checkpoint 为 exactly-once”在架构图上标注端到端 exactly-once。Flink 的 EXACTLY_ONCE checkpoint mode 约束的是作业状态更新;端到端记录交付还要求 source 能恢复同一位点、sink 参与 checkpoint 提交,或目标端能用稳定键幂等收敛。应逐个核对 Doris Pipeline connector及其底层 sink 版本的提交协议、失败恢复、目标表模型和配置。换成 Kafka、Iceberg 或自研 sink 后,结论必须重新证明。
Schema 演进是一条控制流,不是一张字段映射表
3.6 Pipeline 用 SchemaOperator 处理 source 发出的 SchemaChangeEvent。Schema evolution 文档定义五种全局行为:exception 捕获到变化即失败;evolve 尝试把所有变化应用到 sink,失败触发全局 failover;try_evolve 尝试应用,不支持时容忍并转换后续数据,但转换不保证无损;lenient 把变化重写成更宽松、尽量不丢数据的目标变化,默认不传播 truncate/drop;ignore 吞掉结构事件,只继续处理未变化字段。
这五种模式不是从“严格”到“宽松”的简单开关。ignore 可能让新增字段永久缺失;try_evolve 的类型转换可能丢字段;lenient 可能通过重命名旧列再新增列保留数据,却改变目标 schema 的形态;evolve 在 sink 不支持某个 DDL 时会让整个作业反复 failover。团队必须按事件类型制定策略,而不是为所有表统一选择最宽松模式。
正向 Schema 实验先增加一个可空列:
ALTER TABLE sales_db.orders
ADD COLUMN channel VARCHAR(24) NULL;
UPDATE sales_db.orders
SET channel='WEB', updated_at=CURRENT_TIMESTAMP(3)
WHERE id=1001;在 exception 模式下,预期 SchemaOperator 抛出异常并停止数据继续进入 sink;证据包括失败的 schema event 类型、源表、目标表与 sink 返回信息。切换到 evolve 前,先在测试环境证明 Doris connector 与目标版本支持 add column,再从受控 savepoint 恢复。成功时目标表出现 channel,id=1001 为 WEB,旧行允许为 null。
反向实验使用高风险变化,例如把 amount DECIMAL(12,2) 缩窄,或者 drop 仍被消费者读取的列。预期 exception/evolve 明确失败;try_evolve/lenient 即使继续运行,也必须检查目标 schema 与值是否被转换,不能把“Job 仍 RUNNING”当成功。truncate/drop 默认不传播的模式下,源目标行数差异是预期保护结果,需要变更流程决定是否手工执行,而不是偷偷补一条目标 DDL。
checkpoint 和 savepoint 分别解决什么
checkpoint 由作业周期性触发,服务于自动故障恢复,通常有保留策略和清理生命周期。savepoint 由操作者触发,服务于版本升级、拓扑变更、集群迁移或回退,生命周期由团队负责。两者都包含状态,但不能互换责任:把 checkpoint 目录永久留作升级档案会受到自动清理影响;把 savepoint 当高频故障恢复会增加操作与存储负担。
第一次提交前已经配置持久状态路径。升级前还要记录 Job ID、Flink/Flink CDC/connector 版本、Pipeline 文件哈希和 savepoint URI。恢复后观察是否存在 unmapped state、算子 UID 变化、状态 serializer 不兼容。允许忽略未映射状态会直接丢弃旧算子状态,只能在明确证明该状态不再需要时使用。回退也要验证旧二进制能读取升级后产生的状态;“保留了 savepoint”并不自动保证双向兼容。
standalone、YARN 与 Kubernetes 怎样选
standalone 最容易理解,适合本地、实验室和已有主机调度体系的小规模共享集群;资源隔离、进程守护和高可用由团队自己负责。YARN 适合已经有 Hadoop 资源池的组织,可用 application 或 session 形态提交,凭证续期、队列隔离和分布式存储访问成为关键。Kubernetes 适合不可变镜像、声明式资源和平台化运维,但 connector JAR、Secret、ServiceAccount、checkpoint 对象存储权限与 Pod 重建必须一起设计。
Flink CDC 3.6 分别提供 standalone、YARN和Kubernetes提交入口。Kubernetes Operator 入口要求自定义镜像包含 distribution 与 connector,并采用 parent-first 类加载;官方 3.6 文档仍明确原生 application mode 不受支持。选择标准不是“哪种更云原生”,而是状态存储是否可靠、升级路径是否可演练、依赖能否不可变交付、网络与凭证是否能最小化、平台团队是否能接管故障。
Kubernetes 镜像中应固化 Flink、Flink CDC 与 connector JAR,不在 init container 每次从公网拉依赖。Secret 以文件或环境注入时,注意 Flink Web UI、异常堆栈和 Pod spec 是否会暴露明文。NetworkPolicy 只允许 source、sink、checkpoint 存储和必要控制面;ServiceAccount 不应同时具备修改集群工作负载与读取全部业务 Secret 的权限。
项目接入时先定义 sink 幂等契约
每条 Pipeline 都应附一份目标写入契约:目标表模型、业务主键、delete 语义、版本比较字段、重复事件处理、乱序处理、Schema 事件支持集、事务/批次标识、失败后悬挂事务清理,以及对账不变量。没有主键的源表需要特别处理;快照 chunk 与日志阶段都可能依赖分片键,sink 也无法凭空知道两条 append 是否代表同一行。
关系型或主键分析表常用 UPSERT(key, value, source_version),仅当 source_version 不旧于现值才覆盖;delete 记录同样携带可比较版本,防止旧 update 在 delete 后重放导致“复活”。Kafka sink 可用稳定 key 保证同 key 分区顺序,但消费者和 compact 策略仍决定最终状态。对象存储/湖表要区分 append 原始层与 merge-on-read/upsert 表,不能拿原始文件数直接等同业务行数。
项目发布时,Pipeline YAML、connector lock 清单、source/sink DDL、权限申请、savepoint 操作和对账 SQL 应由同一变更关联。新增捕获表不仅要改 tables,还要验证权限、快照容量、server-id 范围、目标路由和 Schema 策略。scan.newly-added-table.enabled 需要从 savepoint/checkpoint 恢复才会对新增表执行重新快照与日志接续;启用前要按 3.6 MySQL connector 文档核对行为,不能期待运行中改正则就自动完整补表。
用状态和指标定位故障,而不是反复重提作业
作业长期停在快照阶段时,先看每表的 isSnapshotting、numTablesRemaining、split 处理数量与最慢 chunk,再检查源库慢查询、主键分布和连接池。所有 reader 都忙但 checkpoint 超时,可能是大 chunk 或背压;reader 空闲而日志 reader 追不上,则看 binlog 产生速率、网络与下游吞吐。
checkpoint 连续失败要按边界分型。source 无法快照状态时看位点与状态大小;barrier alignment 时间升高说明上下游速度不均;SchemaOperator RPC 超时通常表示下游结构变更迟迟未完成;sink pre-commit 卡住则检查目标事务、批次和连接。直接把 timeout 调大只会延迟失败暴露,并可能扩大 source 日志保留。
从 savepoint 恢复报状态不兼容时,比较 operator UID、拓扑、connector 版本与状态 serializer,不要先删除 savepoint重跑 initial。重新全量会改变源库压力、目标重复和切流水位。若源 binlog/WAL 已经过期,即使 savepoint 可读也无法继续,应冻结目标接管,重新建立全量基线并对账。
目标端数据不一致时,把故障窗口内的 source 位点、checkpoint ID、sink 批次/事务 ID 和业务 key 放在一起。相同 key、相同源版本重复通常指向至少一次重放未收敛;key 不同但业务唯一键相同,说明 key 设计错误;一段连续源位点完全缺失,则优先检查 startup mode、日志保留与错误恢复,而不是用去重解释。
容量、权限、成本和团队治理
容量要分别预算 snapshot 和 streaming。snapshot 的瓶颈通常是源库随机/顺序读、并发连接、chunk 倾斜、网络与目标批量写;streaming 的瓶颈是峰值日志速率、反序列化、Schema 协调、checkpoint 状态上传和 sink 提交。pipeline.parallelism 是全局默认,不保证每个 source/sink 都能等比例扩展;MySQL 日志流有天然顺序边界,目标端并发过高还可能形成小批次与事务争用。
状态成本近似由活跃 split、未完成事务、算子状态、checkpoint 频率和保留份数共同决定。对象存储容量看起来充足,也要预算 API 请求、上传带宽、跨区流量与恢复读取时间。团队应记录每条 Pipeline 的源端读放大、日志保留、TaskManager CPU/内存、状态字节、目标写放大和共享集群资源份额,按趋势发现失控。
权限分成四套:source 只读与复制权限、sink 建表/变更/写入权限、checkpoint 存储读写权限、Flink 控制面提交与取消权限。开发者可以查看脱敏指标和日志,不应默认读取所有连接密码;平台管理员能操作集群,也不必拥有业务表全读。Schema 自动演进若开启,sink 账号的 DDL 权限尤其敏感,应限制目标库并记录审计。
Pipeline 日志、dead-letter 数据、checkpoint/savepoint 和失败样本都可能包含敏感字段。transform 中做掩码只能保护 transform 之后的数据,source 内存、状态和错误堆栈仍可能出现原值。敏感数据策略要覆盖加密、访问审计、保留期、删除请求与备份副本;调试时使用合成主键和脱敏值,不把真实行贴进群聊。
团队登记每条作业的 source owner、sink owner、平台 owner、版本组合、状态路径、RPO/RTO、Schema 策略、幂等键、容量预算、告警和退出条件。告警至少覆盖作业状态、连续 checkpoint 失败、checkpoint 时长/大小趋势、source lag、snapshot 剩余、Schema 失败、sink 错误与对账差异。阈值由基线和恢复目标决定,而不是复制另一条作业的数字。
受控回滚与退出不能从删除状态开始
配置或 connector 升级失败时,优先停止新版本,从升级前 savepoint 使用旧镜像、旧 Pipeline 和旧 connector 集合恢复。回退前确认新版本是否已对目标 Schema 做不可逆变化;如果已经 drop/缩窄列,仅恢复 Flink 状态不能恢复目标结构与数据。必要时先执行目标端兼容修复,再恢复旧作业。
迁移结束时先冻结或确认源写入,等待 source lag 与 sink pending transaction 清零,记录最终源位点、最后成功 checkpoint、目标对账结果与业务切流状态。保留一个明确的回切窗口,窗口结束后再停止作业、删除 checkpoint/savepoint、回收 connector JAR 镜像、撤销账号和清理目标临时表。直接删除状态目录会让正在运行的作业在下一次恢复时失去锚点。
清理本地 standalone 实验可执行:
# 先在 Web UI/REST 记录 Job ID 与最终 checkpoint,再取消作业
"$FLINK_HOME/bin/flink" cancel <JOB_ID>
"$FLINK_HOME/bin/stop-cluster.sh"
# 仅删除这次合成实验确认不再需要的状态目录
rm -rf /var/lib/flink/checkpoints/<LAB_JOB_ID>
rm -rf /var/lib/flink/savepoints/<LAB_SAVEPOINT_ID>共享存储的删除应由生命周期策略和审批执行,不能把示例 rm -rf 带到生产。退出完成的判据是源端日志保留责任解除、目标接管证据完整、状态与敏感数据按策略销毁、账号和网络权限回收、作业与告警从资产台账下线。Flink CDC 的价值不只是把数据搬得快,而是让全量、增量、Schema、恢复和外部提交都能被同一条证据链解释。
