news 2026/9/23 11:08:02

NiFi 1.21.0 实现 MySQL 单表增量同步最佳实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
NiFi 1.21.0 实现 MySQL 单表增量同步最佳实践

简介:本资源是一套基于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 个核心组件连成一条流:
GenerateFlowFileExecuteSQLSplitJsonJoltTransformJSONUpdateAttributePutDatabaseRecordLogAttribute

关键设计逻辑:GenerateFlowFile每 30 秒触发一次,生成一个含start_dateend_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 URLjdbc: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 Usernifi_reader不要用 root!创建专用账号:CREATE USER 'nifi_reader'@'%' IDENTIFIED BY 'StrongPass123!'; GRANT SELECT ON mydb.order_info TO 'nifi_reader'@'%'; FLUSH PRIVILEGES;
PasswordStrongPass123!明文填,NiFi 会自动加密存储
Driver Class Namecom.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-addressmy.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 Typesjava.sql.Types.TIMESTAMP,java.sql.Types.TIMESTAMP

  • Query Parameters${start_date},${end_date}

  • Max Wait Time30 sec(防止慢查询拖垮整个流)

关键点在于start_dateend_date的来源——它们由上游GenerateFlowFile的动态属性注入。双击GenerateFlowFileProperties→ 找到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:空值、日期、主键冲突的三重防护

双击PutDatabaseRecordProperties

  • Destination Tableorder_info(目标库表名,必须与源表结构一致)
  • Schema Access StrategyInherit Record Schema(模板已内置 Avro Schema,描述每个字段类型)
  • Statement TypeINSERT不是 UPDATE!增量同步靠主键唯一约束拦截重复,而非ON DUPLICATE KEY UPDATE
  • Allow Missing Columns:✅ 勾选(当源表新增字段而目标表未同步时,不报错跳过)
  • Translate Field Names:✅ 勾选(自动把 JSON 字段user_id映射为数据库列user_id,无需手动配)

最关键的Advanced Settings展开后:

属性作用
Null String LiteralNULL当 JSON 中字段值为字符串"NULL"时,写入数据库NULL,而非字符串'NULL'
Default Values{"status":"pending","amount":0.0}对源数据中缺失的statusamount字段,填默认值(防NOT NULL约束失败)
Batch Size1000每批提交 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

原因GenerateFlowFileRun Schedule设为0 sec(即无限触发),导致 CPU 被占满,NiFi 调度器无法分配线程给下游处理器
解决:双击GenerateFlowFileSettings → Scheduling Strategy→ 改为Timer drivenRun 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

同时在PutDatabaseRecordBatch 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 次后进死信”,需手动配置:

  1. PutDatabaseRecord添加失败关系
    双击PutDatabaseRecordSettings → Relationships→ 勾选failure(默认不勾选)→ 点击Apply

  2. 创建死信队列(Dead Letter Queue)
    拖入一个PutFile处理器 → 命名为DLQ-MySQL-Failures→ 配置Directory/data/nifi/dlq/→ 将PutDatabaseRecordfailure关系线连到它。

  3. 添加重试逻辑(推荐 3 次)
    PutDatabaseRecordDLQ-MySQL-Failures之间插入RetryWithBackoff处理器(需提前安装:下载nifi-retry-bundle-1.21.0.nar放入./lib/)→ 配置Max Retries=3Backoff Interval=10 sec

  4. 对接 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。若全通过,再切回增量模式——这是我对“能跑”和“敢上线”划的生死线。

希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 11:05:04

146、MLIR的Bfloat16与FP8等低精度格式支持

MLIR的Bfloat16与FP8等低精度格式支持 从一次诡异的精度损失调试说起 去年做AI推理引擎时,遇到一个让人抓狂的bug:模型在GPU上跑得好好的,换到某款AI加速芯片上,精度直接崩了。排查了三天,最后发现是MLIR的TypeConverter在把F32转成Bfloat16时,悄悄把某些中间结果的精度…

作者头像 李华
网站建设 2026/9/23 11:04:32

STM32外部中断EXTI原理与实战避坑指南

1. 什么是STM32外部中断&#xff1f;它到底解决什么实际问题&#xff1f;你刚拿到一块STM32最小系统板&#xff0c;接好按键、光耦或霍尔传感器&#xff0c;想一按就立刻响应——结果发现主循环里轮询检测按键状态&#xff0c;CPU一直在空转&#xff1b;或者用定时器不断扫描&a…

作者头像 李华
网站建设 2026/9/23 11:02:44

iCollections 9.6.3:macOS桌面管理的三维革命

1. 工具定位与核心价值解析iCollections 9.6.3 是 macOS 平台上老牌桌面管理工具的最新迭代版本。作为一款专注解决"桌面杂乱症"的专业软件&#xff0c;它通过虚拟分区、智能规则和可视化标签三大核心功能&#xff0c;将传统文件夹的二维管理升级为三维空间管理。我在…

作者头像 李华
网站建设 2026/9/23 11:00:10

实体店如何通过直播带货突破流量困局

1. 实体店与直播间的流量困局最近走访了几家线下实体店&#xff0c;发现一个有趣现象&#xff1a;工作日下午3点&#xff0c;店里一个顾客都没有&#xff0c;但店主却对着手机滔滔不绝。走近一看&#xff0c;原来是在做直播带货。这种"线下冷清、线上热闹"的场景&…

作者头像 李华
网站建设 2026/9/23 10:59:55

分布式风电场LVRT仿真建模与稳定性优化

1. 项目背景与核心价值风电作为清洁能源的重要组成部分&#xff0c;其并网稳定性直接关系到电力系统的安全运行。在实际运行中&#xff0c;电网电压骤降&#xff08;低电压穿越工况&#xff09;是风电机组面临的最严峻挑战之一。去年某省电网的故障统计显示&#xff0c;约37%的…

作者头像 李华
网站建设 2026/9/23 10:58:41

YOLO打火机检测:X光安检小目标识别实战指南

简介&#xff1a;本资源是面向计算机视觉与安防检测领域的YOLO目标检测实践数据集&#xff0c;专为机场X光安检场景中打火机识别任务设计&#xff0c;适用于深度学习初学者、算法工程师及安检系统研发人员。数据集包含2119个真实安检场景图像样本&#xff0c;其中706张JPG格式原…

作者头像 李华