十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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 管理所有状态全存在本地 SQLiteNiFi 自带启动即用日志里每条记录都有 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-jdkUbuntu 默认装的是 JDK 17推荐用 SDKMANcurl -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 兼容性有 regressionwget 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 握手 bugwget 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/nifibin/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.port8081。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属性的 FlowFileExecuteSQL用这两个属性拼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?useSSLfalseserverTimezoneAsia/ShanghaiallowPublicKeyRetrievaltruezeroDateTimeBehaviorconvertToNull必须加serverTimezoneAsia/Shanghai否则DATETIME字段读出来是 UTC 时间zeroDateTimeBehaviorconvertToNull防止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-address在my.cnf中设为0.0.0.0或注释掉。3.2 配置 ExecuteSQL让 WHERE 条件真正按日期增量双击ExecuteSQL处理器 →Configure → PropertiesSQL select querySELECT 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.TIMESTAMPQuery Parameters${start_date},${end_date}Max Wait Time30 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→PropertiesDestination Tableorder_info目标库表名必须与源表结构一致Schema Access StrategyInherit Record Schema模板已内置 Avro Schema描述每个字段类型Statement TypeINSERT不是 UPDATE增量同步靠主键唯一约束拦截重复而非ON DUPLICATE KEY UPDATEAllow Missing Columns✅ 勾选当源表新增字段而目标表未同步时不报错跳过Translate Field Names✅ 勾选自动把 JSON 字段user_id映射为数据库列user_id无需手动配最关键的Advanced Settings展开后属性值作用Null String LiteralNULL当 JSON 中字段值为字符串NULL时写入数据库NULL而非字符串NULLDefault Values{status:pending,amount:0.0}对源数据中缺失的status或amount字段填默认值防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 ?条件永远为FALSEMySQL 中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 中追加zeroDateTimeBehaviorconvertToNull并在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 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.SSS4.5 现象LogAttribute显示flowfile.uuidxxx但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 Retries3Backoff Interval10 sec。对接 Prometheus 监控编辑conf/nifi.properties取消注释nifi.metrics.reporter.prometheus.enabledtrue nifi.metrics.reporter.prometheus.port9092启动后访问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。若全通过再切回增量模式——这是我对“能跑”和“敢上线”划的生死线。希望帮到你。本文还有配套的精品资源点击获取
返回列表