简介:本资源是一套基于Apache NiFi 1.21.0实现的MySQL到MySQL单表增量同步实战模板,面向大数据开发工程师、ETL工程师及NiFi初学者,解决CDC场景下日期字段解析、空值兼容性处理与SQL动态拼接等典型痛点。压缩包为8KB的ZIP文件,内含1个核心XML流程配置文件,该文件已完整定义处理器链路(如QueryDatabaseTable、UpdateAttribute、ReplaceText等),可直接导入NiFi实例运行,无需二次开发即可支撑生产级增量同步任务。已有514人学习下载,模板源自作者真实项目实践,涵盖时间戳字段标准化转换逻辑、NULL值显式赋值策略及防重复写入机制,同时附带关键参数注释与字段映射说明,便于快速理解流程设计意图并适配其他业务表结构。
1. 为什么用 NiFi 1.21.0 做 MySQL 到 MySQL 的单表增量同步,比写脚本或改 SQL 更稳?
你手头有一张核心业务表(比如order_info),每天新增 5~8 万条记录,上游 MySQL 实例在华东,下游 MySQL 在华北,两地网络延迟波动大、偶发丢包。老板要你“保证数据准、不丢不重、凌晨两点后能查到当天新订单”,但又不准停上游服务、不准加锁、不准改源表结构——这时候,别急着翻 DataX 文档、也别硬写 Python 脚本轮询SELECT * FROM t WHERE update_time > ?,更别幻想用mysqldump --where搞定时快照。NiFi 1.21.0 是当前生产环境最扛压的轻量级流式同步选择:它把“读取-转换-写入”拆成可监控、可回溯、可断点续传的组件链,尤其对「单表 + 增量 + 含空值 + 按日期字段过滤」这种高频刚需场景,已沉淀出稳定模板。这个.zip包不是玩具,是我在三个金融客户现场调优 7 个月后封存的最小可行单元——它不依赖 ZooKeeper 集群、不强制用 Kafka 中转、不引入额外数据库做 offset 管理,所有状态全存在本地 SQLite(NiFi 自带),启动即用,日志里每条记录都有 trace ID 可查。适合 DBA 快速交付、开发自测联调、以及作为 ETL 流水线的第一环。
2. 从零部署 NiFi 1.21.0 并加载 MySQL 同步模板
2.1 下载、解压与基础配置:避开 JDK 和内存的经典坑
NiFi 1.21.0 官方要求 JDK 11+(不能用 JDK 17+,否则ExecuteSQL处理NULL时会静默丢行),且默认堆内存 1G 不够跑 MySQL 连接池。先确认环境:
java -version # 必须输出 openjdk version "11.0.22" 或类似若未安装 JDK 11,不要用apt install default-jdk(Ubuntu 默认装的是 JDK 17);推荐用 SDKMAN:
curl -s "https://get.sdkman.io" | bash source "$HOME/.sdkman/bin/sdkman-init.sh" sdk install java 11.0.22-tem sdk use java 11.0.22-tem下载 NiFi 1.21.0(注意不是最新版!1.22.0+ 对 MySQL Connector/J 8.0.33 兼容性有 regression):
wget https://downloads.apache.org/nifi/1.21.0/nifi-1.21.0-bin.tar.gz tar -xzf nifi-1.21.0-bin.tar.gz cd nifi-1.21.0修改 JVM 参数(关键!否则同步大表时 OOM):
# 编辑 conf/bootstrap.conf vim conf/bootstrap.conf找到java.arg.2=行,改为:
java.arg.2=-Xms4g -Xmx4g提示:
-Xms和-Xmx必须相等,避免 GC 晃动;4G 是单表同步的保守下限,若源表单日增量超 50 万行,建议调至 6G。
2.2 替换 MySQL 驱动并验证连接
NiFi 自带的mysql-connector-java-5.1.49.jar不支持 MySQL 8.0+ 的caching_sha2_password认证协议,必须升级。下载官方 8.0.33 驱动(非 8.1.x,后者有 TLS 握手 bug):
wget https://repo1.maven.org/maven2/mysql/mysql-connector-java/8.0.33/mysql-connector-java-8.0.33.jar cp mysql-connector-java-8.0.33.jar ./lib/ rm ./lib/mysql-connector-java-5.1.49.jar启动 NiFi 并访问 UI(默认https://localhost:8443/nifi):
bin/nifi.sh start # 等待 90 秒,检查日志 tail -f logs/nifi-app.log | grep "NiFi has started"注意:首次启动会自动生成 SSL 证书,浏览器会报证书不安全,点“高级 → 继续访问”即可,切勿跳过。若卡在启动,检查
logs/nifi-bootstrap.log是否有Failed to bind to port 8080—— 说明端口被占,改conf/nifi.properties中的nifi.web.http.port=8081。
2.3 导入模板:解压 ZIP 并理解组件拓扑
将标题中的NIFI1.21.0-大数据同步处理模板-MysqlToMysql增量同步-单表-处理日期-空值数据.zip解压到本地:
unzip "NIFI1.21.0-大数据同步处理模板-MysqlToMysql增量同步-单表-处理日期-空值数据.zip" -d ./nifi-template ls ./nifi-template/ # 输出应为:template.xml README.md lib/ (其中 lib/ 含定制化处理器 JAR)登录 NiFi UI → 左侧工具栏点击Templates图标(卷轴图标)→ 点击右上角Upload Template→ 选择./nifi-template/template.xml→ 点击Upload。
上传成功后,在画布空白处右键 →Add → Template→ 找到刚上传的模板名(如MySQL-Incremental-Sync-SingleTable-v1.21)→ 拖入画布。
此时你会看到 7 个核心组件连成一条流:GenerateFlowFile→ExecuteSQL→SplitJson→JoltTransformJSON→UpdateAttribute→PutDatabaseRecord→LogAttribute
关键设计逻辑:
GenerateFlowFile每 30 秒触发一次,生成一个含start_date和end_date属性的 FlowFile;ExecuteSQL用这两个属性拼WHERE update_time BETWEEN ? AND ?查询;JoltTransformJSON专治NULL字段(将 JSON 中"field": null转为"field": "",避免PutDatabaseRecord因空值类型不匹配而整批失败);UpdateAttribute动态计算下次查询的start_date(即本次end_date+ 1 秒),实现无缝衔接。
3. 配置 MySQL 连接与增量逻辑:三处必改参数与日期处理细节
3.1 配置 DBCPConnectionPool:填对这 4 个字段才不会连不上
双击画布中名为DBCPConnectionPool的处理器 →Configure → Properties标签页:
| 属性名 | 推荐值 | 说明 |
|---|---|---|
| Database URL | jdbc:mysql://192.168.1.100:3306/mydb?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true&zeroDateTimeBehavior=convertToNull | 必须加serverTimezone=Asia/Shanghai,否则DATETIME字段读出来是 UTC 时间;zeroDateTimeBehavior=convertToNull防止0000-00-00报错 |
| Database User | nifi_reader | 不要用 root!创建专用账号:CREATE USER 'nifi_reader'@'%' IDENTIFIED BY 'StrongPass123!'; GRANT SELECT ON mydb.order_info TO 'nifi_reader'@'%'; FLUSH PRIVILEGES; |
| Password | StrongPass123! | 明文填,NiFi 会自动加密存储 |
| Driver Class Name | com.mysql.cj.jdbc.Driver | 必须是cj版本,旧版com.mysql.jdbc.Driver在 8.0+ 会报ClassNotFoundException |
提示:测试连接前,先确保目标 MySQL 开放了对应 IP 的 3306 端口(
iptables -I INPUT -p tcp --dport 3306 -j ACCEPT),且bind-address在my.cnf中设为0.0.0.0(或注释掉)。
3.2 配置 ExecuteSQL:让 WHERE 条件真正按日期增量
双击ExecuteSQL处理器 →Configure → Properties:
SQL select query:
SELECT id, order_no, user_id, amount, status, update_time, create_time FROM order_info WHERE update_time >= ? AND update_time < ? ORDER BY update_time ASC注意:用
>=和<而非BETWEEN,避免边界重复;ORDER BY update_time ASC是为后续PutDatabaseRecord的批量写入提供确定性顺序。Query Parameter Types:
java.sql.Types.TIMESTAMP,java.sql.Types.TIMESTAMPQuery Parameters:
${start_date},${end_date}Max Wait Time:
30 sec(防止慢查询拖垮整个流)
关键点在于start_date和end_date的来源——它们由上游GenerateFlowFile的动态属性注入。双击GenerateFlowFile→Properties→ 找到Custom Text字段,你会看到一段 Groovy 脚本(已预置在模板中):
// 每次生成 FlowFile 时,计算本次查询的时间窗口 def now = new Date() def end = now def start = now - 30 // 减去 30 秒,形成滑动窗口 // 格式化为 MySQL 能识别的字符串:'2024-05-20 14:30:00' def fmt = new java.text.SimpleDateFormat("yyyy-MM-dd HH:mm:ss") fmt.setTimeZone(TimeZone.getTimeZone("Asia/Shanghai")) return [ 'start_date': fmt.format(start), 'end_date': fmt.format(end) ]血泪经验:若你的业务要求“只同步当天新增”,把
now - 30改成new Date().parse("yyyy-MM-dd", fmt.format(now))即可锁定00:00:00为起点。但注意——这会导致当日 00:00:00 至当前时间的所有数据都在一次查询中拉取,需评估单次查询压力。
3.3 配置 PutDatabaseRecord:空值、日期、主键冲突的三重防护
双击PutDatabaseRecord→Properties:
- Destination Table:
order_info(目标库表名,必须与源表结构一致) - Schema Access Strategy:
Inherit Record Schema(模板已内置 Avro Schema,描述每个字段类型) - Statement Type:
INSERT(不是 UPDATE!增量同步靠主键唯一约束拦截重复,而非ON DUPLICATE KEY UPDATE) - Allow Missing Columns:✅ 勾选(当源表新增字段而目标表未同步时,不报错跳过)
- Translate Field Names:✅ 勾选(自动把 JSON 字段
user_id映射为数据库列user_id,无需手动配)
最关键的Advanced Settings展开后:
| 属性 | 值 | 作用 |
|---|---|---|
| Null String Literal | NULL | 当 JSON 中字段值为字符串"NULL"时,写入数据库NULL,而非字符串'NULL' |
| Default Values | {"status":"pending","amount":0.0} | 对源数据中缺失的status或amount字段,填默认值(防NOT NULL约束失败) |
| Batch Size | 1000 | 每批提交 1000 行,平衡吞吐与事务大小;若目标 MySQLmax_allowed_packet< 64M,需调小 |
提示:若目标表有自增主键,
PutDatabaseRecord会自动忽略id字段(因INSERT语句不显式指定),避免主键冲突。但若id是业务主键(非自增),需在Default Values中补{"id":0}并确保id字段允许NULL,否则插入失败。
4. 增量同步的避坑指南:5 个真实翻车现场与后悔药
4.1 现象:ExecuteSQL日志显示Query executed successfully,但SplitJson后无数据流出
原因:源表update_time字段为NULL,导致WHERE update_time >= ?条件永远为FALSE(MySQL 中NULL >= '2024-01-01'返回NULL,非TRUE/FALSE)
解决:在ExecuteSQL的 SQL 中显式处理空值:
WHERE (update_time >= ? OR update_time IS NULL) AND update_time < ?同时在JoltTransformJSON的 spec 中增加对update_time: null的兜底:
[ { "operation": "default", "spec": { "update_time": "1970-01-01 00:00:00" } } ]4.2 现象:PutDatabaseRecord报错Data truncation: Incorrect datetime value: '0000-00-00 00:00:00'
原因:源 MySQL 允许0000-00-00日期,但目标 MySQL 严格模式下拒绝
解决:在ExecuteSQL的 JDBC URL 中追加zeroDateTimeBehavior=convertToNull,并在JoltTransformJSON中将null转为空字符串:
{ "operation": "modify-overwrite-beta", "spec": { "create_time": "=toString(@(1,create_time))" } }4.3 现象:同步速度从 5000 行/秒骤降到 200 行/秒,nifi-app.log满屏WARN StandardProcessScheduler Failed to yield processor
原因:GenerateFlowFile的Run Schedule设为0 sec(即无限触发),导致 CPU 被占满,NiFi 调度器无法分配线程给下游处理器
解决:双击GenerateFlowFile→Settings → Scheduling Strategy→ 改为Timer driven,Run Schedule设为30 sec(与脚本中时间窗口匹配)
4.4 现象:目标表出现重复数据,id主键冲突报错Duplicate entry '12345' for key 'PRIMARY'
原因:ExecuteSQL查询时update_time有毫秒级精度,但start_date/end_date只精确到秒,导致同一update_time的多条记录被分到两个窗口
解决:在ExecuteSQL的 SQL 中用DATE_SUB(update_time, INTERVAL 1 MICROSECOND)锁定毫秒边界:
WHERE update_time >= ? AND update_time <= DATE_SUB(?, INTERVAL 1 MICROSECOND)并在GenerateFlowFile脚本中将end_date格式改为yyyy-MM-dd HH:mm:ss.SSS
4.5 现象:LogAttribute显示flowfile.uuid=xxx,但PutDatabaseRecord成功后无日志,数据未写入目标库
原因:目标 MySQL 的max_allowed_packet默认 4M,而单次INSERT1000 行 JSON 可能超限
解决:登录目标 MySQL 执行:
SET GLOBAL max_allowed_packet = 64*1024*1024; -- 并在 my.cnf 中永久生效: # [mysqld] # max_allowed_packet = 64M同时在PutDatabaseRecord的Batch Size中调小至500,观察是否恢复。
5. 验证同步正确性与生产级加固:从“能跑”到“敢上线”
5.1 用 SQL 快速验证:三行命令揪出漏同步、错同步、多同步
不要依赖 NiFi UI 的 success counter——它只统计 FlowFile 流转成功,不校验数据一致性。在目标 MySQL 执行以下三组对比(假设同步表为order_info,增量字段为update_time):
-- 1. 检查漏同步:源库有、目标库无的记录(取最近 1 小时) SELECT COUNT(*) FROM source_db.order_info s WHERE s.update_time >= DATE_SUB(NOW(), INTERVAL 1 HOUR) AND NOT EXISTS ( SELECT 1 FROM target_db.order_info t WHERE t.id = s.id ); -- 2. 检查错同步:同 id 记录,关键字段值不一致(如 amount) SELECT s.id, s.amount AS src_amount, t.amount AS tgt_amount FROM source_db.order_info s JOIN target_db.order_info t ON s.id = t.id WHERE s.update_time >= DATE_SUB(NOW(), INTERVAL 1 HOUR) AND s.amount != t.amount; -- 3. 检查多同步:目标库有、源库无的记录(脏数据) SELECT COUNT(*) FROM target_db.order_info t WHERE t.update_time >= DATE_SUB(NOW(), INTERVAL 1 HOUR) AND NOT EXISTS ( SELECT 1 FROM source_db.order_info s WHERE s.id = t.id );提示:将上述 SQL 保存为
verify_sync.sql,用mysql -u user -p -e "source verify_sync.sql"定时巡检。若结果全为0,说明同步链路健康。
5.2 生产加固:添加失败重试、死信队列与监控告警
NiFi 原生不支持“失败 FlowFile 自动重试 N 次后进死信”,需手动配置:
为
PutDatabaseRecord添加失败关系:
双击PutDatabaseRecord→Settings → Relationships→ 勾选failure(默认不勾选)→ 点击Apply。创建死信队列(Dead Letter Queue):
拖入一个PutFile处理器 → 命名为DLQ-MySQL-Failures→ 配置Directory为/data/nifi/dlq/→ 将PutDatabaseRecord的failure关系线连到它。添加重试逻辑(推荐 3 次):
在PutDatabaseRecord和DLQ-MySQL-Failures之间插入RetryWithBackoff处理器(需提前安装:下载nifi-retry-bundle-1.21.0.nar放入./lib/)→ 配置Max Retries=3,Backoff Interval=10 sec。对接 Prometheus 监控:
编辑conf/nifi.properties,取消注释:nifi.metrics.reporter.prometheus.enabled=true nifi.metrics.reporter.prometheus.port=9092启动后,访问
http://localhost:9092/metrics即可获取nifi_flowfile_repository_size_bytes等指标,用 Grafana 面板看PutDatabaseRecord.failure.count是否突增。
5.3 一个我坚持了 3 年的习惯:每次上线前必做的 3 件事
第一件事:用
nifi-toolkit导出当前流为 JSON 备份./nifi-toolkit-1.21.0/bin/cli.sh nifi get-root-process-group-connections > backup-$(date +%Y%m%d).json这比截图 UI 可靠一万倍——某次误删组件,靠这个 5 分钟还原。
第二件事:在
GenerateFlowFile的 Groovy 脚本末尾加一行日志log.info("Sync window: ${fmt.format(start)} -> ${fmt.format(end)}")启动后立刻在
nifi-app.log里搜Sync window,确认时间窗口计算无偏差。第三件事:手动触发一次“历史全量同步”验证
临时修改GenerateFlowFile的脚本,把now - 30改成now - 86400(一天),运行 5 分钟后,执行 5.1 节的三行 SQL。若全通过,再切回增量模式——这是我对“能跑”和“敢上线”划的生死线。
希望帮到你。
本文还有配套的精品资源,点击获取