
简介这是基于Apache NiFi 1.21.0的MySQL到MySQL增量同步流程模板专为需要做单表CDC实时同步的大数据开发、ETL工程师准备。模板由作者在实际项目中提炼而成导入NiFi后即可直接运行省去从零搭建数据同步流程的重复工作。模板核心实现了基于CDC的增量数据捕获、SQL动态拼接并针对日期类型字段与空值数据做了专门处理能够有效避免同步过程中的格式转换与空指针问题。整个资源包为zip格式内部仅包含1个xml流程定义文件体积约8KB轻量且便于导入和二次编辑。目前已有513人学习下载适用于正在搭建MySQL增量同步管道、或希望参考NiFi模板设计思路的中高级开发者。通过阅读xml中的Processor连线与参数配置可以快速理解增量同步的实现细节并迁移到自己的业务场景中。1. 从一张几亿行的表说起为什么 MYSQL 增量同步比你想的更依赖模板接手过 MySQL 大数据量同步的人应该都有过这种体验第一次全量同步跑完了以为万事大吉结果第二天业务方说“昨天新加的数据没过来”。你查了半天发现不是 NIFI 没跑而是上次同步的断点根本没记对。更隐蔽的是源表里某些字段允许 NULL同步过去后目标表却变成了空字符串消费者的报表在 SUM 和 AVG 时直接算错。NiFi 1.21.0 里那个名为MysqlToMysql增量同步-单表-处理日期-空值数据的模板本质上就是把这套流程固化成了一份可重复使用的资产用 QueryDatabaseTable 记录增量断点high-water mark、用 ExecuteSQL 或查询处理器处理日期边界、再用 PutDatabaseRecord 批量落库。它适合两类人一类是刚接触 NIFI、不想从零理解处理器之间复杂关系的初学者另一类是每天要接几十张表的平台工程师——他们要的不是“能跑”而是“参数改得少、断点不出错、NULL 不被吞”。这篇文章就按模板内部的真实数据流来拆解从增量标记怎么存到日期字段怎么处理再到空值在前端到目标库之间如何保持原样最后落到参数调优和验证手段。全程基于 NIFI 1.21.0 的组件行为不涉及某份不存在的官方文档。2. 模板组件拓扑增量同步的四个核心环节2.1 为什么是 QueryDatabaseTable 而不是 ExecuteSQL 做增量抽取NiFi 里能查 MySQL 的处理器有不少但QueryDatabaseTable是这个场景下最稳的起点。它的核心机制是你指定一个日期或时间戳列或自增 ID 列NIFI 会把上一次查询返回的最大值记录下来下一次轮询时自动带上WHERE update_time 上次最大值这个条件。这个“最大值”被存储为maxvalue属性默认存成 FlowFile 的 attribute模板里通常会引出一个 State 存储来持久化。用 ExecuteSQL 自己拼增量条件也能实现但问题在于断点需要自己维护要么写进一张状态表要么靠变量过期时间硬抗。QueryDatabaseTable的另一个优势是它天然支持分页通过Max Rows Per FlowFile参数控制每次拉多少行避免一张几千万行的表一次性压进内存。这个参数在模板里一般设的是 5000 到 10000具体取值下面有专节讲。注意QueryDatabaseTable 要求你指定的增量列必须是索引列或主键列否则每次轮询都会触发全表扫描几行数据的小表无所谓大表会直接把源库拖死。2.2 单表场景下的流程编排从 GenerateFlowFile 到 PutDatabaseRecord这个模板叫“单表”意味着不需要像多表同步那样用 RouteOnAttribute 按表名分流。它的 FlowFile 路径非常直GenerateFlowFile → QueryDatabaseTable → UpdateAttribute → PutDatabaseRecord或者QueryDatabaseTable → ConvertJSONToSQL → PutSQL。前者适合源表和目标表结构基本一致的情况后者适合需要 SQL 层转换的情况。以最常见的版本为例QueryDatabaseTable输出的每个 FlowFile 是一个 Avro 数据文件里面是增量查询的完整结果集。这里容易踩的坑是很多人以为 QueryDatabaseTable 输出的是一行一个 FlowFile实际上它默认会把整个结果集打包成 1 个 FlowFile后续处理器如果再对这个文件做拆分数据才会分块。模板中通常不拆而是直接让 PutDatabaseRecord 按批次写入。2.3 State Map 的推荐配置你不想在重启后重复灌数据NIFI 1.21.0 的 QueryDatabaseTable 支持两种状态存储方式Local State 和 Zookeeper。单节点验证时用 Local 就够但集群部署时如果不指定 Zookeeper 状态提供者多个节点各自维护断点会出现重复抽取或漏抽的问题。模板里推荐的做法是在nifi.properties里单独配置一个状态提供者专门给同步任务用避免和其他流程的状态混在一起nifi.state.management.provider.clusterzk-provider nifi.state.management.provider.locallocal-provider在 QueryDatabaseTable 的配置界面中State Manager Service下拉框中选中这个服务。很多人忽略这个选项默认的In-Memory State只在当前 JVM 生命周期内有效一旦 NIFI 重启断点归零全量重跑一次——如果你的数据是从 Kafka 转发到 MySQL 的重跑意味着重复写入库业务方会收到两倍的数据。2.4 最小可用拓扑的完整配置参考为了让你能直接照抄下面给出一套最小可用配置。这套配置在单机 NIFI 1.21.0 上验证过源库和目标库都是 MySQL 8.0增量列使用update_time{ processors: [ { name: QueryDatabaseTable, type: org.apache.nifi.processors.standard.QueryDatabaseTable, properties: { Database Connection Pooling Service: MySQL_Connection_Pool, Table Name: orders, Columns to Return: id,user_id,amount,status,update_time, Additional WHERE Clause: , Initial Max Value: 1970-01-01 00:00:00, Max Rows Per FlowFile: 5000, Maximum Value Column: update_time } }, { name: UpdateAttribute, type: org.apache.nifi.processors.attributes.UpdateAttribute, properties: { set_target_table: orders_target } }, { name: PutDatabaseRecord, type: org.apache.nifi.processors.standard.PutDatabaseRecord, properties: { Record Reader: AvroReader, Statement Type: INSERT, Table Name: orders_target } } ] }Columns to Return务必显式列出字段不要用*。一旦源表加了列Avro 文件里会多出新字段而目标表的 INSERT 语句若仍然按旧列数拼PutDatabaseRecord 会直接 D 到 failure 关系。显式指定列名还有一个好处你可以把不需要同步的列比如内部标记字段直接屏蔽掉。3. 处理日期增量都是成也日期败也日期3.1 日期列选型DATETIME 和 TIMESTAMP 的差异化处理MySQL 里DATETIME和TIMESTAMP在 NIFI 增量同步中的行为截然不同。DATETIME存储时不含时区信息NIFI 读出来是什么就是什么直接作为增量条件没有歧义。而TIMESTAMP在存储和读取时会经过 MySQL 会话时区转换如果你在 NIFI 的连接池里没有指定serverTimezone参数JVM 默认时区和 MySQL 时区不一致时读出来的 maxvalue 比实际值早 8 小时或晚 8 小时下一轮增量就会重复抽取。推荐在连接池的连接串里这样写jdbc:mysql://127.0.0.1:3306/source_db?useSSLfalseserverTimezoneAsia/ShanghaiuseCursorFetchtrueuseCursorFetch这个参数容易被忽略MySQL 默认在流式读取时是一次性把结果集拉进内存JVM 堆大表同步时 OOM 就是这么来的。开启游标模式后NIFI 会按fetchSize分块读取内存压力会显著降低。fetchSize在 QueryDatabaseTable 里对应的属性是Max Rows Per FlowFile的底层实现但如果你在 JDBC URL 里没有开启useCursorFetchtrue这个参数对 MySQL 是不生效的。3.2 增量边界用还是以及处理同秒数据的姿势QueryDatabaseTable 生成的 SQL 默认是SELECT id, user_id, amount, status, update_time FROM orders WHERE update_time ? ORDER BY update_time这里的?是上一次记录的最大update_time。注意是严格大于这意味着如果源表在 14:00:00.500 写入了一条记录而上一轮 maxvalue 恰好也是 14:00:00.500同一毫秒有多条数据那么这一条会被下一轮漏掉。解决方式有几种。第一种是在写入源表时保证update_time精度足够比如精确到微秒且业务上不会出现同一精度内的并发写第二种是允许轻微的重复把条件改成然后在目标端用主键去重第三种最优雅——重新生成一个批次内唯一的 maxvalue让下一轮从这个时间点开始比如MAX(update_time) INTERVAL 1 MICROSECOND。但 QueryDatabaseTable 不支持自定义这个表达式你需要用自定义 SQL 的流程。3.3 日期字段的格式统一从 Avro 到目标表的 TIMESTAMP 类型映射源表是DATETIME到了 NIFI 的 Avro 记录里会变成逻辑类型logicalType读取后是字符串表示的时间。此时如果目标表字段也是 DATETIMEPutDatabaseRecord 能自动转换但如果目标表设计成了VARCHAR(20)你需要在前面的处理器里把格式固定下来// UpdateRecord 或 JoltTransform 里可用的表达式 ${field.value:format(yyyy-MM-dd HH:mm:ss)}在 NIFI 里更可控的做法是用ConvertAvroToJSON把数据转成 JSON以yyyy-MM-dd HH:mm:ss格式再通过UpdateRecord的value.provider统一替换日期字段的值。这比依赖数据库端的隐式转换要稳定得多——NIFI 的 Avro 逻辑类型在 1.21.0 中对datetime的输出格式是 ISO-8601如2025-06-01T13:45:30Z直接插入 MySQL 的 DATETIME 会被解析为无效值。3.4 模板里没有告诉你但你必须知道的“次日 0 点”问题你用update_time 2025-06-01 00:00:00做了增量条件目标表对账时发现 6 月 2 日下午 2 点跑的任务把 6 月 2 日 0 点 0 分 0 秒到 2 点之间的数据全部重插了一次。原因在于某个上游任务在 6 月 1 日 23:59:59 写入了数据但事务提交时间晚于 0 点实际的update_time是 6 月 2 日 00:00:00 之前或之后的微妙差别。这类问题本质上是“业务时间”和“系统时间”不一致。模板里如果带了“处理日期”字段通常指的是让你把业务日期如biz_date当作增量依据而不是物理update_time。但物理同步场景下业务日期无法覆盖延迟写入的场景。我的惯例做法是给增量列加一个冗余的“写入时间”字段由业务代码在 INSERT 时无条件写当前时间且设置为NOT NULL保证每一行都有准确的物理时间可供游标追踪。4. 空值数据NULL 的传播路径与处理策略4.1 NULL 在数据库和 NIFI 之间的两层表示MySQL 和 NIFI 对 NULL 的处理有两个层级的差异。第一层是 JDBC 层面ResultSet.getObject()是能返回null的但很多连接池配置里带了zeroDateTimeBehaviorconvertToNull这会让合法的0000-00-00日期变成 NULL造成目标表写入时意外出现 NULL 而不是报错。第二层是 Avro 层面NIFI 生成的 Avro 文件里NULL 值对应的逻辑类型是[null,string]的 union 或type:[null,long]这样的空联合类型。很多解析器对这类复杂 Avro 类型支持不全于是把 NULL 当成了空字符串来处理。验证办法在 NIFI 里用QueryDatabaseTable跑一个 WHERE 条件为column IS NULL的查询然后右键 View FlowFile看数据是否还能看出这是 NULL 而不是空串。如果你在 View 里看到的形如amount:说明某个环节把 NULL 转换成了空字符串如果在 View 里看到amount:null则是正常的。4.2 三种处理姿势让 NULL 保持 NULL、转成默认值、在 SQL 层拦截模板里针对“空值数据”的处理通常有这几种走向。第一种保持 NULL 原样到目标库。做法很简单不要在任何处理器中调用ReplaceText或UpdateAttribute对可疑字段做空值替换同时确认目标表字段没有NOT NULL约束。AvroWriter 和 PutDatabaseRecord 天然支持 NULL 写入 NULL。第二种将 NULL 替换为业务默认值比如将数值型字段的 NULL 替换为 0将字符型字段的 NULL 替换为空字符串。推荐在 SQL 层做掉而不是在 NIFI 做因为这样可以避免 Avro 类型被改动SELECT id, user_id, COALESCE(amount, 0) AS amount, COALESCE(status, ) AS status, update_time FROM orders WHERE update_time ?这种情况下QueryDatabaseTable 不适用你需要改用 ExecuteSQL并使用自定义查询。ExecuteSQL 没有内置的 maxvalue 机制所以你得自己把上一次的同步点位存到一个 NIFI 变量里。这个变量可以在 GenerateFlowFile 时写入 attribute再由 ExecuteSQL 用${last_max_update_time}引用。第三种数据进目标库前主动过滤掉含 NULL 的记录或把 NULL 字段所在行整条丢弃。这适合那些不允许空值的下游表。用RouteOnAttribute配合表达式语言判断${field_name:isNull()}必需字段含 NULL 时将其路由到 failure 而不是丢弃这样排错时有迹可循。4.3 大表上的 NULL 处理不要指望逐行脚本在 NiFi 里写 JavaScript 逐行处理 NULL 是最容易的无底洞。1.21.0 中的 UpdateRecord、JoltTransform 都是流式的但 ExecuteScript 是逐 FlowFile 加载到 JVM 再操作几百万行的 FlowFile 会直接把堆内存挤爆。在这类模板中处理大规模空值我的建议优先级是SQL 层 COALESCE 数据库端视图 NIFI 的 ReplaceText。能不下推给计算引擎的尽量不下推。4.4 practical模板中空值属性参数的推荐配置如果你的模板在设计上带了一个“空值处理”的配置属性一般会见到Null Value Default这样的参数。在 UpdateRecord 处理器里配置如下{ replacementValueStrategy: USE_PROVIDED_VALUE, replacementValue: ${field_name:isNull():ifElse(0, field_name)} }这段表达式的含义是判断field_name是否为空如果为空则替换为字符串0否则保留原值。注意这里0是字符串如果目标字段是整数类型需要在 UpdateRecord 的读取器中明确字段类型。否则 NIFI 会以字符串0写入JDBC 驱动勉强能转但如果你后面做了类型强转或 CAST这里就可能成为性能瓶颈。还有一种情况是数据源里不仅有空值还有null这个字符串字面量。一模一样的时间点写入的字段值可能是null四个字符这不会在isNull()中命中但业务端会把它当文本来处理。处理这个时要多加一个判断${field_name:equals(null):or(${field_name:isNull()})}5. 部署与参数调优从能跑到跑稳的距离就差这几个旋钮5.1 连接池的 Validation Query 与隔离级别连接池配置里有一项很多人不改Validation Query。它每次从连接池借出连接时会执行一次用来防止 MySQL 服务端把空闲连接断开MySQL 默认wait_timeout是 8 小时。如果连接被服务端断掉而客户端不知道NIFI 会报Communications link failure的错然后整个流程进入 retry 循环。模板中推荐把 DBCPConnectionPool 的处理时间设得合理一些初始化连接数2最大连接数10最大等待时间500 millis。注意500 millis不是让你等待数据库响应的时间而是连接管理队列的等待时间。当所有连接被占满时等待超过这个时间的请求直接失败而不是无限阻塞。MySQL 事务隔离级别建议用默认的READ_COMMITTED或REPEATABLE_READ。对于增量查询用 REPEATABLE_READ 会在极端情况下产生间隙锁和幻读风险但 NIFI 的 QueryDatabaseTable 是 SELECT 只读操作不会锁表。重点在于目标库侧的 PutDatabaseRecord如果你的目标库也开启了 REPEATABLE_READ大批量插入死锁的概率会更高建议单独给这个连接池设置TRANSACTION_READ_COMMITTED。5.2 批量参数Max Rows Per FlowFile 与事务大小的取舍Max Rows Per FlowFile设 5000目标库是普通 MySQL 单实例这个值是安全的。如果你把它提到 100000一次性插入的事务会变得很长binlog 文件和 undo log 都会暴涨主从延迟也可能一下飙到几十秒。反过来说设得太小比如 500会导致 PutDatabaseRecord 频繁提交每次提交都有 IO 消耗整体吞吐反而更低。1.21.0 中 QueryDatabaseTable 的每个流程运行可以产生多个 FlowFile当总行数超过 Max Rows 时它会把数据切分为多个 FlowFile每个 FlowFile 都带有相同的maxvalue属性。这一点在集群部署时尤其要注意多个节点并行处理同一批 FlowFile 时如果下游没有加分布式锁可能出现目标端重复插入。建议的参数起点是中转场景参数位置参数名推荐值说明QueryDatabaseTableMax Rows Per FlowFile5000单 FlowFile 数据量控制内存水位QueryDatabaseTableMax Wait Time10 seconds超过该时间无新数据则退出当前轮询PutDatabaseRecordBatch Size1000单事务内写入的行数PutDatabaseRecordStatement TypeINSERT如需幂等可改为 UPDATEDBCPConnectionPoolMax Total10控制源库与目标库的连接使用上限5.3 集群环境下的调度与并发粒度模板若部署在 NIFI 集群上3 节点默认情况下所有节点都可能执行同一个处理器。对于同步类任务需要强制Primary Node Only为true否则每个节点各自维护 State断点就会互相覆盖。这个选项在 QueryDatabaseTable 的 Scheduling 标签页里Permit Scheduling 默认是ALL NODES设成Primary Node Only后同一时刻只有一个主节点跑增量查询其它节点则作为备份等待。如果单个表的增量数据量日均超过 100 万行而且目标表写入瓶颈不大那么主节点单线程抽取的吞吐可能不够。这种情况我更建议把 NIFI 的同步任务只做 CDC 采集和数据落盘输出到 Parquet 文件再由后续的批量任务做分析。不要试图在一个 NIFI 流程里把抽取、清洗、加载全部扛完。5.4 失败重试与幂等落库不要让重试变成重放模板到了目标库写入环节要注意失败重放的幂等性。QueryDatabaseTable 轮询成功后数据被拆成多个 FlowFile如果中间某个 FlowFile 落库失败而你已经把断点推进到了 maxvalue下一次轮询时失败的那部分数据就没有机会再被查出来。解决思路有两个。其一在 PutDatabaseRecord 里设置Update Keys让语句变成 INSERT ... ON DUPLICATE KEY UPDATE这样即使重复插入也只是覆盖旧值其二把失败路由保存下来通过PutFile存成bad-rows文件人工或后续任务补插。前者适用于数据天然有唯一键的表后者适用于无主键的大宽表。模板一般选择方案一因为成本最低但对目标表没有主键或唯一索引的表无能为力。6. 更通用的做法把模板改成可配置的通用同步骨架走到这一步的读者通常已经不只是想跑一个表了。把这张单表模板改造成多表通用骨架只需要动三处。第一处是参数抽离。把 NIFI 模板里的连接串、表名、增量列、只读列、目标表名从编辑器的写死值变成Parameter Context里的变量。比如schema.table、incremental_column、target_table填到参数上下文里这样复制一份模板再改参数就行不需要在处理器界面上到处找。第二处是把 QueryDatabaseTable 的 FlowFile 增加一个target_table属性然后在下游用 RouteOnAttribute 分流到不同的 PutDatabaseRecord。这里要注意PutDatabaseRecord 里表名不能从 attribute 里动态取必须在处理器里配好。所以实用做法是RouteOnAttribute → SetFlowFileAttribute到结果里然后不同分支用不同的 PutDatabaseRecord。第三处是空值策略从“写死在 SQL 里”改成“NIFI 侧按字段做映射”。用 UpdateRecord 的/regex路径定位字段或者用 Jolt 的 modify-overwrite-beta 转换。比如把整条记录的 NULL 值统一替换为0但保留日期字段为空[ { operation: modify-overwrite-beta, spec: { amount: concat((1,amount),), status: toLower((1,status)) } } ]这样模板的可复用性会好很多。你只需要面对新表时重新调整字段映射而不是理解每个处理器的配置含义。增量同步这件事坑大多不在 NIFI 本身而在你对源数据的理解程度日期列会不会回填NULL 值在业务上的含义以及目标表是否真的准备好接收数据。把这三点想清楚了剩下的只是模板里旋钮的微调。本文还有配套的精品资源点击获取