兼容
是对前人努力的尊重
是确保业务平稳过渡的基石
然而
这仅仅是故事的起点

KFS的定位和整体架构

KFS全称是Kingbase FlySync,是电科金仓做的一款异构数据同步软件。它基本上涵盖了数据同步的主流需求。

先说支持的数据源。源端支持Oracle、SQL Server、MySQL/MariaDB、KingbaseES、DB2、Informix、达梦、OceanBase、TDSQL、PolarDB、openGauss、Vastbase、磐维DB、AntDB等主流数据库,以及Kafka、RabbitMQ、RocketMQ等消息队列。目标端除了上述数据库外,还支持ClickHouse、GBase8a、Greenplum、ElasticSearch、MongoDB、Hive、Hadoop等大数据和分析型数据库。加起来总共30多种数据源,覆盖面还是很广的。

整个同步链路是一个流水线结构,分成了四个主要部分:

Pipeline 流水线层
┌────────────┐    ┌───────────┐    ┌─────────────┐    ┌──────────────┐
│ Extractor  │───▶│  Filter   │───▶│ Partitioner │───▶│ Applier × N  │
│  采集变更   │    │ 事件分流   │    │  路由到通道  │    │ 每通道独立   │
│  事件       │    │ 过滤·转换  │    │             │    │ 线程+连接    │
└────────────┘    └───────────┘    └─────────────┘    └──────────────┘

用大白话说就是:

  • Extractor(采集器):从源端数据库的日志里抓变更数据。KFS不走SQL查询,而是直接解析数据库的事务日志(Binlog、Redo、WAL),实时捕获INSERT/UPDATE/DELETE操作。这种方式对源端的侵入性极低,我后面会专门说。
  • Filter(过滤器):对抓到的事件做过滤和转换。比如你只需要同步某几张表,或者需要对某些字段做转换,就在这里处理。还支持正则表达式过滤数据对象,灵活度很高。比如你只想同步order_*开头的表,一行正则就能搞定。
  • Partitioner(分区器):这个是KFS的一个亮点设计。它负责把事件分配到不同的通道里去并行处理。路由策略有好几种,从简单的哈希到负载均衡再到行哈希,后面展开讲。
  • Applier(加载器):每个通道有独立的线程和数据库连接,负责把数据写到目标端数据库。它把统一格式的事件转换成目标端数据库的原生SQL语句,然后执行。

这个设计的好处是什么呢?就是每个环节都可以独立优化、独立扩展。采集慢了可以优化采集策略,加载慢了可以增加通道数,互不干扰。而且KFS支持1对1、1对多、多对1、级联、双向等多种同步拓扑,可以灵活组合。比如你可以一个源端同步到多个目标端(数据分发),也可以多个源端汇聚到一个目标端(数据集中),甚至可以级联同步(A→B→C)。

采集模块:对源端"零干扰"是怎么做到的

对源端数据库来说,KFS基本上是个"旁观者"——你在正常做你的业务,KFS在旁边默默读取日志文件,互不干扰。

当然也支持分离方式部署,就是把采集模块部署在单独的服务器上,进一步降低对生产服务器的影响。对于那些对性能特别敏感的生产系统,分离部署是比较推荐的做法。

我实际测试过,在源端跑了大约2000 TPS的业务负载下,开启KFS同步之后,源端的CPU负载增加不到3%。这个数据跟KFS产品文档里宣传的"低侵入"基本吻合。对于一个实时同步工具来说,这个侵入性已经非常低了。对比以前用过的某些基于触发器方案的同步工具,动不动就让源端CPU多涨20%以上,差距很明显。

KUFL:统一的中间格式

采集模块拿到数据之后,会把它封装成一个统一的中间格式,存储到KUFL文件里。KUFL的全称是Kingbase Unified Format Log,这是KFS自己的跟踪文件格式,用了加密的数据格式。

它的设计目的是什么呢?就是把不同数据库中的数据操作信息,用一种统一的格式来存储。不管你的源端是Oracle的Redo日志还是MySQL的Binlog,到了KUFL里都是统一格式。这样做的好处是:后续的过滤、转换、路由、加载模块就不需要关心源端是什么数据库了,只需要处理统一格式的数据就行。整个系统的模块解耦就做得很好。

# KUFL文件的典型目录结构
$ ls -la /opt/kfs/data/kufl/
-rw-r--r-- 1 flysync flysync  256M  815 10:23 kufl_000001.log
-rw-r--r-- 1 flysync flysync  512M  815 14:47 kufl_000002.log
-rw-r--r-- 1 flysync flysync  128M  815 18:05 kufl_000003.log  # 当前写入中

这种架构还有个好处:当源端或者目标端中断的时候,KUFL文件里已经包含了到中断时刻的最新数据。系统恢复之后直接从这些文件里加载就行,不用重新从源端拉数据。这就为断点续传提供了基础。

另外KUFL文件采用了加密存储,数据传输过程中也支持端到端国密算法加密(SM4)。KFS与数据库连接还支持SSL/TLCP等加密通信协议。在数据安全越来越受重视的今天,特别是金融、政务这些对数据安全性要求很高的行业,这个特性还是很有必要的。KFS的产品代码自主率是100%,已经在厂内完成了产品高风险漏洞处理和开源代码消化/重构。

目标端入库优化:KFS是怎么"追得上"的

前面说了,核心矛盾是"源端几百个连接并行写入,目标端CDC工具逐条回放"。KFS在目标端做了几个很关键的优化,我觉得值得拿出来说一说。这些优化不是孤立的,而是相互配合形成一个完整的性能提升体系。

批量提交——减少网络往返

最基础的优化。以前逐条执行SQL,每条SQL都要走一次网络往返(RTT),假设RTT是1ms,6条SQL就是6ms延迟。KFS改成批量执行,6条SQL打包成一次网络请求,只需要1个RTT。

// 优化前:逐条执行
Statement stmt = conn.createStatement();
stmt.execute("INSERT INTO t VALUES(1, 'a')");  // 1次网络往返
stmt.execute("INSERT INTO t VALUES(2, 'b')");  // 1次网络往返
stmt.execute("INSERT INTO t VALUES(3, 'c')");  // 1次网络往返
// ... 6条SQL = 6次网络RTT,至少6ms延迟

// 优化后:批量执行
PreparedStatement pstmt = conn.prepareStatement("INSERT INTO t VALUES(?, ?)");
pstmt.setInt(1, 1); pstmt.setString(2, "a");
pstmt.addBatch();
pstmt.setInt(1, 2); pstmt.setString(2, "b");
pstmt.addBatch();
pstmt.setInt(1, 3); pstmt.setString(2, "c");
pstmt.addBatch();
pstmt.executeBatch();  // 6条SQL = 1次网络RTT,网络往返减少约5倍

这个优化看起来简单,但对TPS的提升是立竿见影的。特别是在网络延迟比较高的场景下(比如跨广域网同步),效果更明显。你想,如果RTT是10ms(跨地域很正常),6条SQL逐条执行就是60ms,批量执行只要10ms,直接快了6倍。

小事务合并——按表归并减少prepare次数

这个优化比较有意思。假设一个事务里面对10张表轮流做insert,而且是乱序的(表1、表2、表3…表10、表1、表2、表3),如果按源端顺序逐条处理,每切换一张表就要重新做一次prepareStatement。

为什么要关心prepare次数?因为每次prepareStatement,数据库都需要做一次SQL解析——词法分析、语法解析、生成执行计划。这些都是CPU密集型操作。特别是当表很多、事务又频繁切换表的时候,大量的CPU时间被浪费在重复解析同样的SQL模板上。

KFS的做法是:先把同一个表的SQL归到一起,按表排序之后再统一处理。这样一来,每张表只需要prepare一次,后面全部是addBatch操作。

// 优化前:按源端顺序处理
prepare 表1 → execute
prepare 表2 → execute
prepare 表3 → execute
prepare 表4 → execute
prepare 表1 → execute  // 又切回表1了,重新prepare
prepare 表2 → execute  // 又切回表2了,重新prepare
prepare 表3 → execute  // 又切回表3了,重新prepare
共 prepare 7次(表1和表2各出现2次,重复解析)

// 优化后:按表归并
prepare 表1 → addBatch, addBatch  // 表1的所有操作一起执行
prepare 表2 → addBatch, addBatch
prepare 表3 → addBatch, addBatch
prepare 表4 → addBatch
共 prepare 4次(去重后的表数)

// prepare次数:7 → 4,减少约43%

我在测试环境里模拟了一个跨10张表乱序变更的场景,30个事务、每个事务20条记录。开启小事务合并之后,prepare次数从两百多次降到了三十多次,回放耗时降了大概一半。效果还是挺明显的。

而且配合Batch使用,同表连续addBatch还能进一步减少网络往返,形成一个叠加的优化效果。

Statement缓存——避免重复解析SQL

这个也好理解。数据库收到一条SQL,得先做词法分析、语法解析,生成执行计划,这些都是CPU密集型操作。如果同样的SQL模板反复发来,每次都重新解析一遍,那就太浪费了。

KFS在Applier层做了一个SQL模板缓存(大小500),相同模板的SQL只需要解析一次,后面直接复用已经编译好的PreparedStatement。

// SQL模板缓存机制示意
// 使用LinkedHashMap实现LRU缓存,最大500条
Map<String, PreparedStatement> stmtCache = new LinkedHashMap<>(500);

String sqlTemplate = "INSERT INTO " + schemaTable + " VALUES(?, ?, ?)";
PreparedStatement pstmt = stmtCache.get(sqlTemplate);
if (pstmt == null) {
    // 缓存未命中:需要数据库做完整的解析流程
    pstmt = conn.prepareStatement(sqlTemplate);
    // 词法分析 → 语法解析 → 语义检查 → 生成执行计划
    // 以上步骤在prepareStatement时由数据库完成
    stmtCache.put(sqlTemplate, pstmt);
} else {
    // 缓存命中:直接复用,跳过所有解析步骤
    // 只需要重新绑定参数值
}
pstmt.setString(1, val1);
pstmt.setString(2, val2);
pstmt.setDouble(3, val3);
pstmt.addBatch();

// 缓存带来的收益:
// 1. 减少SQL解析(词法+语法分析只做一次)
// 2. 复用执行计划(数据库不用重新优化查询)
// 3. 降低CPU开销(省去语义检查等CPU密集步骤)
// 4. 配合Batch使用(同一个PreparedStatement可连续addBatch)

Buffer Flush策略——吞吐和延迟的平衡

批量提交虽好,但也不能无限攒下去。你想,如果一直攒batch不提交,目标端的数据就一直不可见,同步延迟就上去了。所以需要一个flush策略来平衡吞吐量和数据可见性延迟。

KFS的flush策略支持两种触发方式:定量触发(攒够了N条就刷一次)和定时触发(每隔T毫秒刷一次)。这两个参数需要根据实际场景来调优——如果更在意吞吐量,就把batch size调大;如果更在意实时性,就把flush间隔调短。

多通道并行入库——突破单线程天花板

前面这几个优化都是在单通道内做文章。但单线程再怎么优化,也就一个CPU核心在干活。现在的服务器动辄几十个核心,不用白不用。所以KFS设计了多通道并行入库的机制。

但多线程就带来了一个新问题:怎么保证数据的顺序性和一致性?同一张表的INSERT和UPDATE如果被分到了不同通道,而且UPDATE的先执行了,那数据就乱了。

这就是Partitioner模块要解决的核心问题。KFS的路由策略经历了好几代演进,每一代都是在解决上一代的局限:

第一代:原始单通道——最简单,task直接1:1绑定channel,什么路由逻辑都没有。零开销,但也零并行。

第二代:Hash分区——按shardId做哈希取模分配到不同通道。好处是同一个数据源的事件永远去同一个通道,保证了同源有序。但问题是如果某个shard是热点(比如核心业务库),对应通道就会一直很忙,其他通道闲得慌。

第三代:轮询(RoundRobin)——不管shardId,按事务序号轮流分配。事务数绝对均衡,但忽略了事务大小的差异。一个大事务可能包含几十万行,直接把一个通道塞满。

第四代:负载均衡(LoadBalancing)——实时探测每个通道的队列深度,每次把事件分配给当前行数最少的通道。

// 负载均衡分区
int minQueue = 0;
for (int i = 1; i < queues.length; i++) {
    if (queues[i].size() < queues[minQueue].size()) {
        minQueue = i;
    }
}
// 选队列最短的通道,实时动态均衡
// Q0: 35行  Q1: 5行  Q2: 20行 → 选Q1
// 问题:只看队列深度,不管同一张表是否被分散到不同通道

第五代:表亲和性 + 事务拆分——跨表大事务可以按表拆分成多个子事务,分散到不同通道并行处理;但同时保证同一张表的所有操作始终路由到同一个通道,保证单表有序。

// 事务拆分示意:一个大事务跨3张表
// 拆分前:整个事务堆在通道0
大事务T(表1+表2+表3) → 通道0(其他通道空闲)

// 拆分后:按表拆成3个子事务并行
表1 → 子事务1(fragno=0) → 通道0
表2 → 子事务2(fragno=1) → 通道1
表3 → 子事务3(fragno=2) → 通道2

KFS把这几种策略都封装好了,用户可以根据自己的场景选择最合适的组合,这个灵活度还是很好的。

断点续传:崩溃恢复的一致性保障

做数据同步的,不可能不考虑崩溃恢复的场景。进程挂了、网络断了、数据库重启了,这些迟早都会遇到。如果你的同步工具不支持断点续传,那一次中断就意味着从头来过,这在生产环境中是不可接受的。

KFS的做法是在目标端维护一个叫trep_commit_seqno的表,每个通道一行,记录了该通道已经处理到哪个位点。

-- 多通道断点记录表示意
-- task_id | seqno | fragno | last_frag | wait_seqnos
--    0    |  100  |   2    |   true    | 10|100-2|98-0,99-1|97-3
--    1    |   99  |   1    |   true    | 10|100-2|98-0,99-1|97-3
--    2    |   97  |   3    |   true    | 10|100-2|98-0,99-1|97-3

其中wait_seqnos是一个复合位点,格式是:全局提交计数 | 通道0待处理位点 | 通道1待处理位点 | ...

恢复的时候分三步:

  1. 找到所有通道位点中的最小值(上例中是97),作为安全重发起点
  2. 找到所有通道位点中的最大值(上例中是100),作为已提交的上界
  3. 从最小值开始逐事件重放,对于在97到100之间的事件,检查它是否在wait_seqnos里:
    • 在 → 说明这个事务还没被所有通道处理完 → 重新回放
    • 不在 → 说明已经提交过了 → 跳过

先用KDTS(数据库迁移工具)基于SCN/LSN快照做一次全量数据搬迁,把存量数据搬到目标端。在全量搬迁的同时,KFS的源端增量解析模块就已经开始工作了——它从全量搬迁对应的SCN/LSN号之后开始解析新增的Redo日志,把增量数据缓存在KFS节点本地。等全量搬迁完成后,再启动数据加载,把缓存的增量数据追上去。

全量搬迁 + 增量补偿流程:

时间线 ──────────────────────────────────────▶

源端业务:  ────────持续写入─────────────────────
                  ↑
                  │ 全量搬迁开始(基于SCN/LSN快照)
                  │
全量搬迁:  ══════╡  KDTS搬运存量数据  ╞═════
                  │
增量解析:         ╞═══════════════════════════▶
                   解析SCN之后的增量日志,缓存到本地
                                    
增量加载:                              ╞══════▶
                                        加载缓存的增量数据
                                        直到两端SCN/LSN追平

这个设计的好处是:全量搬迁可以慢慢搬,不用着急,也不用停业务。增量数据先缓存在本地,等全量搬完了再一起追上去。整个过程对用户的应用是完全透明的。

同步拓扑:不只是"一对一"

前面说了KFS支持多种同步拓扑,这里展开聊几个比较实用的。

1对多(数据分发):一个源端同步到多个目标端。典型场景是总部向多个分支机构分发数据。比如安徽省发改委那个案例,合肥中心的数据要实时分发到16个地市,就是一个典型的1对16分发场景。

              ┌──▶ 地市1 KES
              ├──▶ 地市2 KES
省中心 KES ───┤
              ├──▶ 地市15 KES
              └──▶ 地市16 KES

多对1(数据集中):多个源端同步到一个目标端。典型场景是数据汇聚、数据中台。比如解放军总医院的案例,8个医学中心的数据汇聚到一个DWS数仓。

双向同步:两个数据库互相同步。典型场景是双轨并行、双活容灾。比如西安第一人民医院的项目,Oracle和KES之间双向同步,新老系统同时运行,随时可以切换。

Oracle ADG ◀══════════▶ KES集群
  (老系统)    KFS双向同步    (新系统)

级联同步:A同步到B,B再同步到C。适用于多层级数据分发场景。

跨受限网络同步:这个比较特殊。有些场景下源端和目标端之间隔着防火墙、网络隔离设备、甚至单向光闸,网络不是直通的。KFS支持跨正/反向网络隔离设备、跨单向通信光闸、跨物理隔离网络同步。它通过域名解析和网络传输优化,在这些受限网络环境下也能实现数据的可靠同步。广域网传输还有4倍的压缩比,2M带宽就能支撑实时容灾。

我自己比较印象深刻的是湖南移动统一调度平台那个案例。它的源端和目标端数据库网络无法直连,中间隔了承载网。KFS通过域名解析的方式解决了网络不通的问题,400G+存量数据加每天400G+的增量,同步延迟还是在秒级。

支持的数据类型:比你想的多

KFS支持的数据类型远比你想象的要多。除了常规的CHAR、VARCHAR、NUMBER、DATE、INT、FLOAT这些基本类型,它还支持:

  • 大对象类型:BLOB、CLOB、TEXT、XML等。大对象同步一直是个难点,因为体积大、传输慢,但很多业务系统确实需要。
  • GIS数据类型:KFS可以配合Oracle和KingbaseES的GIS能力,无损解析并转换几何存储类型。这对水利、自然资源等行业的国产化替代很重要。如果你用其他同步工具,GIS数据可能只能当字符串处理,空间查询功能就废了。
  • SEQUENCE、函数、存储过程、视图、同义词、索引等数据库对象:KFS不仅同步数据,还能同步数据库对象定义。这在迁移场景下特别有用,你不需要单独去目标端重建这些对象。

另外KFS还支持DDL同步。也就是说你在源端执行一个ALTER TABLE ADD COLUMN,这个变更会自动同步到目标端。这个在双轨并行场景下特别重要——业务系统还在迭代,数据库结构随时可能变,如果DDL不能同步过去,那每次改表结构都要手动在两端操作,太容易出错了。

当然DDL同步也有需要注意的地方:KFS通过NoneCriticalPartitioner确保DDL操作固定在通道0执行,避免和DML并发执行导致元数据冲突。这个设计在前面多通道并行那节已经讲过了。


好了,上篇就先写到这里。主要讲了KFS的整体架构、异构数据同步的难点、数据初始搬迁的方案、同步拓扑的灵活性、目标端入库的几个核心优化(批量提交、小事务合并、Statement缓存、Buffer Flush),以及多通道并行的路由策略演进和断点续传的设计。

下一篇文章,也就是下篇,我会重点讲KFS的数据一致性校验和修复这个核心能力——这也是我觉得KFS跟其他同步工具最大的差异化所在。包括在线校验怎么做的、MD5摘要比对怎么用、自动修复怎么配置、命令行工具怎么用、以及几个真实项目案例里是怎么用的。还会分享一些我在实操中踩过的坑和调优经验。

下篇见。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐