MySQL增量同步模板:NiFi 1.21.0处理日期边界与空值数据
简介面向大数据集成场景的NIFI1.21.0流程模板专门解决MySQL到MySQL的增量数据实时同步问题尤其对日期类型和空值数据进行了针对性处理。模板可直接导入NIFI复用免除从零搭建流程的繁琐工作适合已掌握NIFI基本操作、需要快速落地单表增量同步任务的开发与数据维护人员参考。压缩包内容非常精简共包含1个文件即核心的XML流程模板体积仅8KB将XML导入NIFI后即可看到完整处理流程再按实际项目修改数据库连接、表名及字段映射参数便能投入运行。这一模板在CSDN已有514人学习浏览说明其在实际场景中具备一定的参考与复用价值。模板内部实现了增量CDC数据的抓取逻辑并通过SQL拼接进行动态查询同时覆盖日期格式化、空值判断与默认值填充等易错细节这些处理经验来自作者真实项目实践能帮助读者避开常见同步坑点也可用作学习NIFI数据流设计的入门素材。1. 这个模板解决的是哪种同步痛做过 MySQL 数据同步的工程师多半有过这种经历业务表每天新增几万行全量导入一次越来越慢于是改成只取当天变化的数据结果改了 where 条件之后跨天日期边界漏数NULL 值到了目标库直接报错最后只能手动补数、重跑批。这个标题里的 NiFi 1.21.0 模板干的就是这件事——用 Apache NiFi 的QueryDatabaseTableRecord做基于日期字段的增量抽取再配合PutDatabaseRecord写回目标 MySQL单表单测并把最容易翻车的日期边界和 NULL 值处理固化成模板。不用写代码也能复现适合正在做 MySQL 到 MySQL 增量同步的工程师尤其是用过 DataX 但被脚本维护成本拖累的人。2. 为什么选 NiFi 而不是 DataX 或 binlog增量同步的选型逻辑2.1 增量同步的三条路线与适用边界把数据从 MySQL 挪到另一个 MySQL业界常见的路线其实只有三条。第一条是把目标表按主键最大值来拉SELECT * FROM t WHERE id 上次最大值。这条路线很简单也能保证不重复但它只能处理“只新增不改历史”的数据一旦业务上 update 了老主键的行为就会被漏掉。标题里带上了“处理日期”说明这套模板走的不是纯 ID 轮询而是和日期字段配合的增量判断。第二条是解析 binlog 的 CDC 路线例如 Canal、Debezium。它的优点是近乎实时、能捕获 delete 和 update 的旧值缺点是部署组件多要开 binlog、要配 Kafka、要处理 schema 变更对一个“单表同步”的需求来说太重了。很多团队上一个 CDC 项目光排查 binlog 位点回退就花了两周。第三条就是标题里这个模板走的路基于日期/时间戳字段的增量轮询。每次任务执行时NiFi 记录下这一轮查到的最大日期值下一轮就把WHERE update_time 最大值作为查询条件。你不需要写定时脚本去维护“上一次跑到了哪里”NiFi 自己会把这个状态存在本地。我做同步类任务时选型判断就一句话数据量每天十万行以内、能接受分钟级延迟、表上有可靠更新日期字段就用第三条每天百万行以上、上游表没有日期字段才去上 binlog。这个模板的价值就是把第三条路线里最繁琐的状态维护和类型转换封装成了可视化流程。2.2 NiFi 1.21.0 在 MySQL 同步上的关键能力NiFi 1.21.0 最合适做增量同步的处理器是QueryDatabaseTableRecord。它的行为和旧版QueryDatabaseTable类似但输出的是 Record 格式后续接 JSON、Avro、SQL 转换都很顺。它有几个对增量同步特别有用的设计。Maximum-value Columns参数用来指定一个或多个“增量临界列”。NiFi 执行查询时会自动生成形如SELECT id, ... FROM table WHERE update_time 2024-06-01 12:00:00的 SQL并把查询结果集里这一列的最大值保存到 state 中作为下一轮的查询起点。整个状态机的判断逻辑由处理器内部完成你不用在流程里去拼 SQL。其次是“结果拆包”。一次查询返回的数据行数可能很大QueryDatabaseTableRecord会根据Max Results和Fetch Size参数把数据拆成多个 FlowFile避免一个庞大的 JSON 文件把内存撑爆。这一点对于“大数据同步处理”场景很重要单表几百万行时拆成每个一万行的 FlowFile下游一边读一边写不会出现 OOM。目标端的关键处理器是PutDatabaseRecord。它可以读取 Record 格式的 FlowFile然后把每条记录映射成 INSERT 或 UPDATE 语句。配合目标表的主键或唯一索引可以实现“有则更新、无则插入”的写回效果。这个特性用在增量同步里很省心——如果哪一轮因为日期边界重叠导致同一行被拉取了两次目标端不会因为主键冲突直接报错退批。2.3 单表 MysqlToMysql 模板的标准流程拓扑这套单表模板的流程拓扑不复杂典型链路是源 MySQL → QueryDatabaseTableRecord → UpdateRecord → PutDatabaseRecord → 目标 MySQLQueryDatabaseTableRecord负责抽取UpdateRecord负责处理日期格式和 NULL 空值PutDatabaseRecord负责写入目标库。三条连接分别走success和failure失败的数据单独落到一个日志文件或告警队列里不阻塞主链路。有人会问为什么中间夹一个UpdateRecord不能直接源端到目标端写吗原因就在标题后半段——“处理日期、空值数据”。源 MySQL 的 DATETIME 字段读进 NiFi 后在 Record 里可能被理解为字符串也可能被理解为带时区的 Timestamp而源库某个字段为 NULL目标表列恰好设置为 NOT NULL。这些情况不在中间层做一次字段改写写入动作就是要反复试错的。所以中间层的存在不是为了炫技而是把脏数据挡在写库之前。下面两章我会把该层的配置讲透。3. 手动建立这个流程从 NiFi 1.21.0 启动到 MySQL 双端配置3.1 准备环境启动 NiFi 1.21.0 与两个 MySQL 库NiFi 1.21.0 是 2023 年发布的版本要求运行在 JDK 8 或 JDK 11 上。本地验证时直接解压二进制包启动即可。tar -xzf nifi-1.21.0-bin.tar.gz cd nifi-1.21.0/bin ./nifi.sh start # 等待约 30 秒然后检查状态 ./nifi.sh status启动后默认 Web 端口是 8080浏览器访问http://服务器IP:8080/nifi即可。NIFI 默认是单用户登录模式吗不是。1.21.0 默认启用了 HTTP 访问但如果你下载的包启用了鉴权会要求打开/conf/nifi.properties把nifi.security.user.login.protocol改为非鉴权模式或用 OIDC。本地测试最常见做法是临时改回非鉴权生产环境再用独立账号体系。源库和目标库的连接串我建议提前在 MySQL 侧确认一下权限。同步账号不能只给 SELECTPutDatabaseRecord对目标库需要 INSERT/UPDATE 权限。最小授权如下CREATE USER sync_user% IDENTIFIED BY SyncPassw0rd; GRANT SELECT, SHOW VIEW ON source_db.* TO sync_user%; GRANT INSERT, UPDATE, DELETE ON target_db.* TO sync_user%; FLUSH PRIVILEGES;3.2 源端表结构设计日期字段与空值字段的模拟为了把整个流程跑通我们需要在源库里建一张业务表并在目标库建一张结构接近但未必完全一致的表。-- 源库表 CREATE TABLE source_db.user_business ( id BIGINT PRIMARY KEY AUTO_INCREMENT, user_code VARCHAR(64) NOT NULL, user_name VARCHAR(128), status TINYINT, remark TEXT, last_update_time DATETIME NOT NULL, create_date DATE, INDEX idx_update (last_update_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; -- 目标库表结构与源表基本一致但允许 remark 为空 CREATE TABLE target_db.user_business_sync ( id BIGINT PRIMARY KEY, user_code VARCHAR(64), user_name VARCHAR(128), status TINYINT, remark TEXT, last_update_time DATETIME, create_date DATE ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;注意两点。第一源表的last_update_time必需有索引否则每次增量查询都会全表扫描。第二目标表的remark字段可空因为源端可能写入 NULL。真正生产表单里这一步的字段映射表需要仔细比对不能用 Navicat 直接同步表结构了事。插入几行测试数据要故意制造边界情况INSERT INTO source_db.user_business (user_code, user_name, status, remark, last_update_time, create_date) VALUES (A001, 张三, 1, 正常数据, 2024-06-01 09:00:00, 2024-06-01), (A002, 李四, 0, NULL, 2024-06-02 10:30:00, 2024-06-02), (A003, 王五, 1, , 2024-06-03 23:59:59, 2024-06-03);第二行是典型的“空值数据”陷阱remark为 NULL直接 INSERT 到非空字段会报错。第三行是空字符串和 NULL 不是一回事。处理方式也不同后面会讲。3.3 控制器服务与核心处理器配置进入 NiFi 界面以后先把左侧组件区的Controller Services打开新建两个 DBCPConnectionPool分别指向源库和目标库。这样处理器里就不需要重复填连接串了。源库连接池的关键参数如下参数名配置值Database Driver Location/opt/nifi/lib/mysql-connector-j-8.0.33.jarDatabase Driver Class Namecom.mysql.cj.jdbc.DriverDatabase Connection URLjdbc:mysql://source-host:3306/source_db?useSSLfalseallowPublicKeyRetrievaltrueserverTimezoneAsia/ShanghaiDatabase Usersync_userPasswordSyncPassw0rd目标库连接池同源库只需把 URL 换成 target-host 和 target_db。然后拖入三个处理器按上一章的拓扑连线。第一个是QueryDatabaseTableRecord这是流程的增量抽取端。配置如下# 处理器关键配置 Record Reader → JsonTreeReader或 AvroReader JDBC Connection Pool → source-dbcp Table Name → user_business Maximum-value Columns → last_update_time Initial Value → 2024-06-01 00:00:00 Fetch Size → 5000 Max Results → -1这里Initial Value指定第一次抽取的起始日期。如果不填NiFi 默认从表的最小值开始拉相当于全量。生产环境首次上线时建议先手工指定一个起始时间做完整量快照再改为增量启动。第二个是关键的处理中间层UpdateRecord用于处理 NULL 和日期格式。它会把读到的 Record 字段进行改写。配置方法如下Record Reader → JsonTreeReader Record Writer → JsonRecordSetWriter Replacement Value Strategy → Record Path Value # 添加一个属性改写 /remark → ${ field(remark):isEmpty() ? : field(remark) }${ field(remark) }表达式对 NULL 和空字符串的判断通过:isEmpty()统一收敛。这里的逻辑是NULL 或空串都改写成空字符串避免目标端非空校验失败。如果业务上需要保留 NULL 而不是空串表达式反过来写即可。第三个是PutDatabaseRecord负责把数据写进目标库JDBC Connection Pool → target-dbcp Record Reader → JsonTreeReader Statement Type → INSERT Table Name → user_business_sync如果希望“更新已存在主键”把 Statement Type 改成UPDATE并在目标表建立主键约束。更进一步源表有 delete 操作需要传递时思路就完全不同了需要改用 CDC这里不展开。3.4 把模板导出成 zip 备份的做法NiFi 里的流程模板可以右键流程空白处选择 “Export”导出一个.xml模板文件。Controller Service 配置不会自动打入模板所以常规做法是把关键 DBCP 连接串写在模板描述里。这个标题里给出的 zip 包推测就是流程模板 XML 加说明文档的打包产物。你自己如果搭建出了可用流程建议维护同样的 zip 形式——对新环境部署导入模板后只需改两个连接池的连接串就能跑这比让运维对照文档重新拖流程图可靠得多。导入路径是模板图标 → Import → 选择 XML。导入后模板会出现在组件列表里直接拖到画布就会生成一串处理器和关系连线配置项需要重新填。这不是 NiFi 的 bug而是设计如此模板只保存拓扑和属性默认值不保存连接凭据。4. 增量日期与空值参数的配置矩阵4.1 日期增量字段Maximum-value 的真正含义Maximum-value Columns这个名字很容易误导人它并不是求表的某列最大值而是指定“把哪一列作为增量水位”。NiFi 会执行完查询后把该列返回结果的最大值写入 state。选择增量列有两条硬规则。第一列必须是可排序且单调递增的常见选择是UPDATE_TIME、CREATE_TIME、自增 ID。第二如果选择了日期列要保证表内的日期精度统一。源的 DATETIME 是秒级精度NiFi 记录水位也按秒存如果某个业务写入端手动改成了毫秒精度那么水位会短暂停在旧的秒值上下一轮可能重复拉取秒级边界的那几行。Initial Value参数决定第一次触发的起点通常填一个比最早业务数据更早的时间。但要注意查询条件默认是严格大于即 value而不是。如果第一次的初始值和某行数据时间完全相等这行会被漏掉。因此初始值最好比真实起点早一秒比如 00:00:00 去拉安全起见写成 前一天的 23:59:59。在 1.21.0 里QueryDatabaseTableRecord的查询行为会通过QueryDatabaseTableRecord的State控制。如果处理器名前缀一样且表名一样它会复用同一份 state。所以你改了表结构或换了目标库之后如果不显式清 state可能沿用旧水位导致数据补不回来。操作路径是处理器右键 → View State → Clear。4.2 空值数据的两条处理策略空值数据分两种形态SQL 的 NULL 和空字符串。这在 MySQL 里是两种完全不同的值。remark TEXT既可能为 NULL也可能为。PutDatabaseRecord 在写入时对 NULL 的处理方式取决于目标表列约束目标列没有 NOT NULL 约束NULL 直写没问题。目标列有 NOT NULLNULL 会导致 SQL 异常整批事务回滚。处理策略常见有两种。第一种是保留语义只做“NULL → 空字符串”的收敛用 UpdateRecord 的表达式改写。第二种是丢掉空值数据即在查询阶段就排除掉关键字段为 NULL 的行这种适合“空值本来就没意义”的场景。但用QueryDatabaseTableRecord没法加自定义 WHERE所以要改用ExecuteSQL自己写 SQL。# ExecuteSQL 方案关键 SQL 示例 SELECT id, user_code, user_name, status, IFNULL(remark, ), DATE_FORMAT(last_update_time, %Y-%m-%d %H:%i:%s) AS last_update_time, create_date FROM user_business WHERE last_update_time ${prev_max_date}用这个方案可以精细控制 NULL 和日期格式但需要自己维护prev_max_date。这也是为什么标题模板选择QueryDatabaseTableRecord——NiFi 帮你维护水位代价是不方便写自定义 WHERE。如果你连日期格式化也想省可以把目标表字段直接设置为 VARCHAR但这会丢掉数据库侧的时间语义不建议。4.3 批次大小、并行度与连接池的边界Fetch Size控制每次从 JDBC ResultSet 拉取的行数默认值 0 代表使用 JDBC 驱动的默认 fetch sizeMySQL 驱动默认是 10非常保守。对同步场景设置 2000 到 5000 是比较合理的太大对内存冲击明显太小则查询次数太多。Max Results代表单个 FlowFile 的最大数据行数。设置成 -1 表示不按行拆一个 FlowFile 装完设置 10000则每搬运一万行就生成一个 FlowFile。后者更利于边拉边写目标库压力均匀。这里建议保持默认 -1等观察 NiFi 堆内存使用后再调成 20000 左右。关于轮询频率QueryDatabaseTableRecord的调度策略可以配置为 CRON 表达式例如每 5 分钟跑一次。调度并发建议保持 1。增量任务不像普通抓取任务并发线程会在 state 读写上产生竞争重复拉取的概率会明显升高。你不需要通过加并发来提升抽取速度真正的吞吐瓶颈在目标库写入侧。DBCP 连接池本身有几个值得调的值Maximum Pool Size 默认 8对单表同步完全够用Validation Query填SELECT 1防止 MySQL 主动断开空闲连接后 NiFi 仍使用失效连接。Max Wait Time保持默认即可同步场景不需要追求过短的等待。5. MySQL 增量同步避坑指南五个最常翻车的现场5.1 日期边界重叠导致同一行数据被重复写入现象第一次跑增量到 10:00第二次跑 10:00 到 11:00第二次的表里出现了两条 id 相同的数据或者目标库直接报主键冲突。原因一种是查询条件用了边界上的数据被拉取了两次另一种是QueryDatabaseTableRecord的状态里记录的边界精度不够比如源库写入时间是2024-06-03 23:59:59.500水位记录成2024-06-03 23:59:59下一轮查询 23:59:59就会同秒的比……其实 MySQL 的 DATETIME 秒级精度不会出现这种情况但一旦业务改成毫秒级问题就会出现。解决统一精度。要么把源表字段改成 DATETIME(3)要么在 NiFi 处理器里主动把水位向下取整。最稳妥的做法是目标表加主键或唯一索引并配合PutDatabaseRecord的 UPDATE 模式。这样即便重复也只是更新不会报错。5.2 NULL 值导致写入报错 “Column ‘remark’ cannot be null”现象同步日志里出现 SQLException目标端表里某行根本没有写进去但 NiFi 界面上流程还停留在 Running 状态看起来像一切正常。原因源端某行 remark 为 NULLQueryDatabaseTableRecord原样把 NULL 带进 RecordPutDatabaseRecord生成的 INSERT 语句里没有这个列的值而目标列是 NOT NULL数据库拒绝执行。解决在 UpdateRecord 里把 NULL 改写成默认值表达式如${ field(remark):isEmpty() ? : field(remark) }。如果业务上要保留 NULL则调整目标表为允许 NULL。二者选其一不能只靠处理器报错后人工补数。另外要注意 NULL 也会影响日期字段——比如create_date为 NULL写入前要么排掉要么赋一个默认日期1970-01-01。5.3 MySQL 8 连接时 SSL 错误或公钥检索失败现象NiFi 报错Communications link failure或者Public Key Retrieval is not allowed。原因MySQL 8 默认要求 SSL 连接同时caching_sha2_password认证插件在首次建连时需要获取服务器公钥。很多人只填了用户名密码没管 JDBC URL 的 SSL 参数。解决在 JDBC URL 上显式加参数jdbc:mysql://host:3306/db?useSSLfalseallowPublicKeyRetrievaltrueserverTimezoneAsia/ShanghaiallowPublicKeyRetrievaltrue只建议在可信内网环境使用公网环境必须换成useSSLtrue并配置证书。另外MySQL 8 驱动要用 8.x 版本别拿 5.1 的老驱动连 8.0 的库那不是版本玄学是协议不兼容。5.4 日期时区漂移同步结果比源库早或晚 8 小时现象同步完成后目标库表里的 last_update_time 比源库小了 8 小时或大了 8 小时日期对不上账。原因NiFi JVM 默认时区和 MySQL 的会话时区不一致。MySQL 返回 DATETIME 不带时区信息但 JDBC 驱动在读取 TIMESTAMP 时会按serverTimezone转换NiFi 侧的 Record Writer 再按 JVM 默认时区序列化两边一加一减就出现了偏移。解决在连接池 URL 中强制指定serverTimezoneAsia/Shanghai并且在启动 NiFi 的nifi-env.sh里加一行NIFI_JVM_OPTS-Duser.timezoneAsia/Shanghai。这样从驱动到 Record 序列化都在同一时区日期字段就不会漂了。5.5 State 状态丢失导致从全量重新抽取现象重启 NiFi 后第一次同步发现目标表数据量暴增像是重新全量同步了一遍。原因QueryDatabaseTableRecord的 state 默认存在本地state目录里。单机重启一般不会丢但如果你把整个 NiFi 目录迁移到另一台机器或者该处理器被删除后重新添加state 就找不回来了。NiFi 在找不到 state 时会按Initial Value重新开始如果你没设 Initial Value它就直接全量拉。解决迁移前先导出 state或者在重新添加处理器时显式指定一个靠近当前实际水位的时间让第一轮增量范围可控。另一种做法是集群部署并启用 ZooKeeper 状态存储让水位不依赖单机磁盘。状态丢失本身不是致命的真正危险的是你不知道它丢了等数据重复了才发现。所以每次重启后先源库和目标库各查一次MAX(last_update_time)对比一下再放量。6. 验证这套方案的三个有效手段6.1 用造数脚本验证增量边界先跑通一次全量确定目标表已有部分数据。然后手工对源库插入一批跨越日期边界的新数据INSERT INTO source_db.user_business (user_code, user_name, status, remark, last_update_time, create_date) VALUES (V001, 验证A, 1, NULL, NOW(), CURDATE()), (V002, 验证B, 0, 正常, NOW(), CURDATE()), (V003, 验证C, 1, NULL, NOW() INTERVAL 1 SECOND, CURDATE());触发 NiFi 跑一轮再到目标库查SELECT COUNT(*) FROM target_db.user_business_sync WHERE user_code LIKE V%。如果只有两条而不是三条说明remark为 NULL 的那行被过滤策略吃掉了不是你期望的结果。如果三条都到了但日期不是今天就要回头查时区配置。6.2 用对账 SQL 验证增量与空值处理验证增量不能只看行数要看数据窗口。源库和目标库执行同一窗口统计SELECT MAX(last_update_time) AS max_time, COUNT(*) AS row_count, SUM(CASE WHEN remark IS NULL OR remark THEN 1 ELSE 0 END) AS null_count FROM source_db.user_business WHERE last_update_time 2024-06-01 00:00:00; -- 目标库执行同样的查询最后对比三个数字对账时把 NULL 也统计进去否则处理策略是丢弃还是收敛为默认值根本看不出来。我在实际工作中习惯把这条 SQL 写成一个 Shell 脚本同步跑完自动比对不等业务来质疑数据。6.3 模拟故障回放清状态重跑验证方案的恢复能力核心动作是本地状态被清掉后能不能安全重跑。在 NiFi 界面找到QueryDatabaseTableRecord右键 → View State → Clear然后置为非调度状态手动触发一次。观察是否还会从全量开始。如果发生了全量但目标表有主键且使用 UPDATE 模式数据也不会炸只是耗时变长。最后一次提醒验证不能只在测试表上做。我吃过亏测试表只有几百行怎么跑都没问题上了正式环境几百万行一跑半小时后内存飙升因为 Max Results 配置的是 -1整个结果集都进了内存。后来我给生产表设置了 Fetch Size 5000、Max Results 20000并让 NiFi 每处理完两个 FlowFile 就停下来看一眼堆内存确认稳定才让定时调度接管。这个习惯帮我拦下了不少问题希望帮到你。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联
返回资讯列表 →