
https://github.com/apache/
点击蓝字
关注我们
导语
在数据集成系统中,一个简单的配置参数往往隐藏着复杂的工程设计。Apache SeaTunnel 7 月 Meetup 专场中,讲师以 JDBC Sink 的 batch_interval_ms 参数作为切入点,深入剖析了数据写入过程中的批处理机制,并将视角从 Connector 层逐步扩展到 SeaTunnel Zeta Engine 的整体架构设计。
讲师介绍 PROFILE
牛志伟
GitHub ID:nzw921rx

在数据同步系统中,很多问题最初看起来都属于某一个具体 Connector,但随着深入分析,往往会发现它们真正涉及的是整个数据处理引擎的运行机制。
这次分享源于 JDBC Sink 中一个看似简单的参数:batch_interval_ms。
这个参数的目标很明确:当 buffer 中的数据等待时间超过设定值时,即使没有达到 batch size,也应该触发一次 flush,从而降低数据同步延迟。
但在实际实现过程中,问题很快超出了 JDBC Connector 本身的范围。
因为它背后涉及的并不是 JDBC 如何执行 SQL,而是一个数据同步引擎需要解决的运行时问题:
更进一步,当 Source 暂时没有数据产生时,系统是否仍然能够按照时间触发 flush,也成为必须考虑的问题。
因此,这次问题最终从 JDBC Sink 的一个参数batch_interval_ms,逐步深入到了 SeaTunnel Zeta Engine 的设计层面,并引出了 STIP-23 中 Engine-Level FlushSignal 的设计。
01 背景:数据延迟
在 JDBC Sink 的批量写入场景中,通常存在两个触发 flush 的条件:batch_size 和 batch_interval_ms。

其中,batch_size 的逻辑比较直接。当数据不断进入 Sink 时,系统将数据放入 buffer,并判断当前缓存数量是否达到设定阈值。如果达到阈值,就执行一次批量写入。
例如:
if (buffer.size >= batchSize) { flush();}这种方式天然适合放在 writeRecord() 中,因为只有新的数据进入时,buffer size 才会发生变化。
但是,batch_interval_ms 的语义完全不同。
它表达的并不是“下一条数据到来时,顺便检查一下距离上次 flush 是否已经超过指定时间”,而是希望实现一种真正基于时间的触发机制:即使当前没有新的数据输入,只要距离上一次 flush 已经达到设定时间,也应该有机会执行 flush。
这两个语义在高吞吐场景下可能差异并不明显,因为数据不断到达,writeRecord() 会持续执行,时间判断也会不断被触发。
但在真实生产环境中,大量任务并不是持续高吞吐运行。
例如 CDC 同步任务、小表同步任务、业务低峰期的数据变化,或者某些 Source 分区暂时没有新的数据产生时,buffer 中可能已经存在等待写入的数据,但由于没有新的 record 到来,系统无法进入下一次 writeRecord() 调用。
此时,batch_interval_ms 想表达的“按时间触发 flush”就无法实现。
这也是为什么这个问题最终不能简单归结为 JDBC Sink 的一个参数问题,而需要从 Engine 的运行机制重新思考。
02 从一个参数开始

最初解决这个问题时,一个自然的想法是在 JDBC Connector 内部增加定时任务。
例如通过 ScheduledExecutorService 创建一个后台线程:
scheduledExecutor.scheduleAtFixedRate(() -> { flush();}, batchIntervalMs, batchIntervalMs, TimeUnit.MILLISECONDS);这种方式确实可以满足最直接的需求。
即使 Source 暂时没有数据产生,后台线程仍然可以按照时间周期调用 flush。
但是,当这个方案放入 SeaTunnel Task 执行模型之后,问题就出现了。
原本的数据处理流程中,writeRecord() 和 flush 都属于 Task 正常执行路径的一部分。但增加后台线程之后,flush 变成了 Connector 自己维护的一条额外执行路径。
Task Thread 中执行:
writeRecord(record)
而 Connector Background Thread 中执行:
flush()
这意味着两个线程可能同时访问同一个 buffer。
与此同时,flush 还可能和 checkpoint、schema evolution、close 等流程同时发生。后台线程中的异常也无法自然传递到 Task 主执行路径,任务 fail、cancel、close 时,还需要额外保证这个定时线程能够正确停止。
如果每个 Sink Connector 都采用类似方式实现定时 flush,那么每个 Connector 都需要重复维护自己的 timer 逻辑和线程生命周期。
因此,问题逐渐暴露出来:
为了实现一个时间参数,Connector 开始承担原本属于运行时系统的职责。
这并不是 Connector 应该解决的问题。
另一种思路是不创建后台线程,而是在 writeRecord() 中增加时间判断。
例如:
public void writeRecord(Row row) { buffer.add(row); if (buffer.size() >= batchSize) { flush(); return; } if (System.currentTimeMillis() - lastFlushTime >= batchIntervalMs) { flush(); }}这种方式避免了额外线程,flush 和 write 始终在同一个执行路径中,异常也可以直接反馈给 Task,生命周期管理更加简单。
但是它无法解决最核心的问题。
因为没有新的数据输入时,writeRecord() 根本不会被调用。
因此,它实际实现的并不是“每隔 5 秒执行一次 flush”,而是“下一条数据到来时,如果发现距离上一次 flush 已经超过 5 秒,就顺便执行一次 flush。”
在数据持续流动的场景中,这种差异并不明显。
但在低吞吐、CDC 间歇变化、小表同步以及 Source 暂时空闲等场景下,buffer 中的数据可能长期等待,而系统无法获得下一次触发机会。
这说明两个方案都无法真正满足 batch_interval_ms的语义。

Connector 内部创建线程虽然能够按照真实时间触发,但会引入并发、异常传播和生命周期管理问题;而在 writeRecord() 中判断时间虽然保持了单线程执行模型,却无法处理 Source 空闲场景。
问题的关键并不在 JDBC 实现,而在于抽象层级。
batch_interval_ms 真正需要的是一种 Engine 能力:
即使没有新的数据输入,Engine 也能够按照时间产生一个控制事件,这个事件不会直接操作 Sink,而是进入正常的数据处理通道,最终由 Sink 在自己的消费线程中完成 flush。

也就是:
Engine Timer → FlushSignal → Sink flushAction
这正是 STIP-23 希望解决的问题。
03 Zeta Task执行模型:
理解FlushSignal进入Engine的基础
在讨论 FlushSignal 为什么需要由 Engine 产生之前,需要先理解 SeaTunnel Zeta Engine 中一个 Task 是如何运行的。因为 FlushSignal 并不是一个独立存在的机制,它需要进入 Zeta 已有的数据处理模型,并沿着 Engine 管理的数据链路向下游传递。
如果不了解 Task 的调度、数据流转以及控制事件处理方式,很难理解为什么 FlushSignal 最终选择作为一种 Engine-Level Signal,而不是由 Connector 自己维护。

在 Zeta Engine 中,一个数据同步任务提交之后,并不是简单地由某一个线程直接执行,而是经过 Engine 的调度和分发过程,将任务拆分并分配到不同 Task 中运行。
Engine 负责整个任务生命周期的管理,包括任务提交、资源协调、Task 创建以及运行状态维护。根据任务的执行计划,Engine 会将不同的数据处理阶段分配到对应 Task 中,由 Task 负责实际的数据处理工作。
从执行模型来看,一个 Task 并不是独立运行的业务逻辑,而是由 Engine 统一管理的执行单元。它包含数据处理过程中的输入、转换和输出逻辑,并按照 Engine 定义的运行模型完成数据传递。
这种设计保证了 Connector 不需要关注任务如何调度,也不需要自己维护执行线程。Connector 只需要实现自身的数据处理逻辑,而 Task 的生命周期、状态管理以及线程模型由 Engine 统一负责。
这也是后续 FlushSignal 设计的重要基础。
因为如果 Connector 自己创建 Timer,并直接调用 flush,本质上就是绕开 Task 的执行模型,建立了一条 Engine 无法感知的额外执行路径。

理解 Task 执行模型之后,还需要进一步了解普通数据 Record 在 Zeta Engine 中是如何传递的。
SeaTunnel 的数据处理链路遵循典型的数据流方向:
Source → Transform → Sink
Source 负责从外部系统读取数据,并将数据转换成 SeaTunnel 内部统一的数据结构。
随后,Record 会进入 Transform 阶段,根据用户配置执行字段映射、过滤、转换等处理逻辑。
处理完成后的数据继续向下游传递,最终进入 Sink,由 Sink Connector 完成目标系统写入。
在这个过程中,Record 并不是直接由 Source 调用 Sink,而是经过 Engine 管理的数据通道完成传递。
这意味着,数据流转过程中的每一个阶段都受到 Engine 的统一管理,包括数据传递、线程模型以及任务状态。
对于 FlushSignal 来说,这一点非常重要。
因为如果 Engine 已经存在一套完整的数据传递链路,那么新的控制事件更合理的方式就是复用这条链路,而不是重新建立一个旁路。
也就是说,FlushSignal 不应该由 Timer 线程直接调用 Sink,而应该像 Record 一样进入 Engine 的数据通道。

除了普通的数据 Record,SeaTunnel Engine 中还存在另一类特殊事件:控制事件。
SchemaChange 就是其中一种典型场景。
在 CDC 或动态 Schema 变化的数据同步任务中,数据结构可能随着源端变化而发生改变。例如新增字段、修改字段类型等情况。
这类变化并不是一条普通的数据记录,而是一种需要同步到下游的结构变化事件。
因此,SchemaChange 需要沿着数据链路进行传递。
它不应该被 Source 单独处理,也不应该绕过 Transform 和 Sink,而应该作为 Engine 数据流中的一种特殊事件,由各个阶段识别并继续向下游传播。
对于 Transform 来说,它通常不需要理解具体的 Schema 变化业务含义,而是保证事件能够继续传递。
最终,Sink 根据自身能力处理 SchemaChange,将源端结构变化同步到目标系统。
这个过程说明,在 SeaTunnel Engine 中,数据流不仅仅承载业务 Record,也能够承载控制类事件。
而 FlushSignal 的设计正是基于这一思路。
如果 SchemaChange 可以作为一种事件在数据链路中传播,那么 FlushSignal 同样可以作为一种 Engine 控制事件进入这条链路。

除了 SchemaChange,Checkpoint Barrier 是 Zeta Engine 中另一类重要控制事件。
在 Exactly-Once 场景下,Checkpoint 用于保证任务状态和数据处理结果的一致性。
Checkpoint 过程中,Engine 会生成 Barrier,并将其注入数据处理链路。
Barrier 会随着数据流向下游传播,在不同阶段之间形成一致的状态边界。
当 Task 收到 Barrier 后,会根据自身状态完成对应处理,例如保存状态、协调 checkpoint 流程,并最终保证整个任务在失败恢复时能够从正确的位置继续执行。
这里体现了一个非常重要的设计思想:
控制事件不需要绕过数据处理链路,而是可以复用 Engine 已有的数据通道。
Barrier 不是通过额外线程直接通知某个组件,而是沿着 Task 数据流进行传播。
这和 FlushSignal 的设计思路高度一致。
FlushSignal 并不是一个特殊的外部调用,而是一种新的控制事件类型。它需要像 SchemaChange 和 Barrier 一样,由 Engine 产生,并通过已有的数据链路传递。
通过 Task 调度与分发、Record 数据流转、SchemaChange 控制事件以及 Checkpoint Barrier 链路这几个部分,可以看到 SeaTunnel Zeta Engine 已经具备了一套完整的数据与事件处理模型。
因此,当 JDBC Sink 遇到 batch_interval_ms 这个问题时,真正需要解决的并不是如何让 JDBC 自己定时调用 flush,而是如何在 Engine 已有模型中增加一种新的控制事件。
这也是为什么最终 STIP-23 选择了 Engine-Level FlushSignal。
它不是给 JDBC 增加一个特殊机制,而是在 SeaTunnel 已有的 Task 和事件模型基础上,增加了一种能够表达“触发 flush”的运行时信号。
04 Flush 的引擎抽象
在 JDBC Sink 的实现过程中,batch_interval_ms 暴露出的核心问题并不是如何执行一次 flush,而是如何让 flush 具备 Engine 级别的调度能力。
传统 Connector 实现通常将 flush 看作 Sink 内部的一次缓存提交动作,当数据量达到阈值时,通过 writeRecord() 判断并触发 flush。但时间触发和数据触发存在本质区别。batch_size 依赖 Record 的持续到来,而 batch_interval_ms 要求即使没有新的数据进入,只要达到时间间隔,也能够触发一次 flush。
因此,flush 不能继续停留在 Connector 方法调用层面,而需要被提升为 Engine 可以管理的一种运行时事件。
STIP-23 的核心思路,就是将 Flush 抽象成 FlushSignal。它不再由 Connector 自己维护 Timer,也不是由独立线程直接调用 Sink 的 flush 方法,而是由 Engine 根据时间产生一个 Signal,并让这个 Signal 沿着 SeaTunnel 原有的数据处理通道进行传播,最终由 Sink 根据自身语义决定如何处理。
在这个模型中,FlushSignal 与 SeaTunnel 中已有的控制事件保持一致。数据链路中不仅存在业务数据 Record,同时也存在 Checkpoint Barrier、SchemaChange 等控制信息。FlushSignal 同样作为一种特殊事件进入 Record 通道,使 Engine 能够统一管理它的生命周期和传播过程。
Flush 从方法调用转变为 Signal,是整个设计变化的起点。
在传统模式下,Flush 通常直接发生在 Sink 内部。例如 JDBC Sink 中,当缓存达到一定大小时,会调用 executeBatch() 将数据写入数据库。这种方式对于基于数据量的触发条件非常适合,因为 Record 的到来会不断推动 buffer 状态变化。
但是,当触发条件变成时间之后,问题开始出现。batch_interval_ms 表达的是“经过一定时间后,即使没有新的数据进入,也应该尝试 flush”。如果仍然依赖 writeRecord() 触发,那么实际效果变成了“下一条数据到来时检查是否超时”,而不是严格意义上的定时 flush。

为了说明不同类型事件在 Engine 中的统一处理方式,我们使用 SeaTunnel 数据通道中的几类数据标识。其中:
R
Ck
Sc
F
R → R → R → Ck → R → R → Ck
表示业务数据和 checkpoint barrier 按照顺序进入处理链路。
当引入 SchemaChange 后,数据流中会出现:
R → R → Sc → R → R → Ck
SchemaChange 作为一种控制事件插入数据流中,但不会改变整体的数据传递模型。
FlushSignal 采用同样的方式:
R → F → Sc → R → R → Ck → R
它不是绕过数据链路直接通知 Sink,而是作为一种新的事件类型进入已有 Record 通道。
因此,FlushSignal 的本质不是新增一个 flush 调用方式,而是让 flush 成为 Engine 可以感知和调度的一种运行时信号。
FlushSignal 的产生由 Engine 负责,而不是 Connector 自己创建后台线程。
SeaTunnel Engine 完整的触发流程:

整个过程从配置开始。
当任务配置:
sink.flush.interval
不为 0 时,表示开启 FlushSignal 功能。如果配置为 0,则表示关闭,不影响已有任务运行。
任务启动过程中,SourceFlowLifeCycle 会完成 Timer 的注册工作。随后,TaskExecutionService 中的 timerFlushWorker 根据固定延迟方式执行周期任务。
这里 Timer 的职责非常明确。它只负责按照配置周期触发任务,并在时间到达时调用 onTimerTick()。Timer 本身并不执行 Sink flush,也不会直接操作 SinkWriter。
当 Timer Tick 发生后,Engine 会调用:
collector.sendFlushSignal(jobId, taskId)
将 FlushSignal 注入 Source 输出。
PPT 特别强调,FlushSignal 的注入需要与 Checkpoint 使用同一把 checkpointLock。
这样设计的原因是,FlushSignal 和 Checkpoint Barrier 都属于影响数据处理状态的控制事件,需要避免两者在注入过程中产生并发冲突。
因此,Engine 不是通过额外线程强制触发 Sink,而是在 Source 侧按照已有任务执行模型生成 FlushSignal,使其进入正常的数据处理流程。
FlushSignal 成功注入之后,会复用 SeaTunnel 原有的 Record 通道进行传播。
Signal 流转的链路为:

其中中间可能经过 Engine 内部的 Queue / Disruptor。
整个过程可以理解为:
在 Source 阶段,FlushSignal 会通过:
sendRecordToNext()
广播给下游所有:
output.received(record)
也就是说,FlushSignal 和普通 Record 使用相同的数据发送机制。
进入 Transform 后,Transform 不需要理解 Flush 的业务含义,也不需要执行任何转换逻辑。
识别 Signal 后直接 collector.collect(record),不进入 transform()。这意味着 FlushSignal 不会参与普通数据转换过程,而是保持原始状态继续向下游传播。
最终,FlushSignal 到达 SinkFlowLifeCycle。
Sink 侧会识别这是一个 Signal,而不是普通业务数据,并进入:
processSignal()
执行后续处理。
因此,FlushSignal 的流转过程遵循了一个重要原则:
它进入已有 Record 通道,但只有 Sink 才真正理解它的业务含义。
FlushSignal 进入数据通道后,还需要处理 Queue 满载情况下的数据流控制问题。
针对不同类型数据记录,可以采取不同的入队策略:

对于普通 DataRecord:
put() / ringBuffer.next()
当 Queue 满时,会等待容量释放,然后继续保留数据。
对于 Barrier:
同样采用:
put() / ringBuffer.next()
当 Queue 满时等待容量,保证 Checkpoint 信息不会丢失。
而 FlushSignal 使用的是不同策略:
offer() / tryPublishEvent()
当 Queue 满时:
立即返回 false
并:
丢弃本次 Signal
也就是说,FlushSignal 并不会像业务数据和 Checkpoint Barrier 一样阻塞数据通道。
这样设计是因为 FlushSignal 本身并不携带业务数据,它代表的是一次 flush 意图,而不是必须保证到达的数据内容。
如果 Queue 已经处于满载状态,说明当前系统正在处理大量数据,此时优先保证业务数据和 checkpoint 正常流转更加重要。
同时,
prepareClose:Signal 直接绕过
以及:
Queue 满:FlushSignalQueueFailureTotal +1
这表示在任务关闭或者异常情况下,FlushSignal 的处理策略也进行了特殊设计,不会因为 Signal 自身阻塞整个数据处理流程。
这种设计保证了 FlushSignal 具备轻量级控制事件的特点,不影响主数据链路稳定性。
FlushSignal 最终到达 SinkFlowLifeCycle 后,由 Sink 负责解释它的实际含义。
在整个流程中,Engine 并不知道 JDBC 的 flush 具体如何实现,也不会直接调用数据库操作。

Engine 只负责将 FlushSignal 送达到 Sink。
SinkFlowLifeCycle 收到 Signal 后进入:
processSignal()
然后根据 Signal 类型执行对应逻辑。
对于普通 Record:
Record → sinkWriter.write()
继续执行正常数据写入。
对于 FlushSignal:
FlushSignal → flushAction
执行 Connector 注册的 flush 行为。
这也是 Engine 和 Connector 职责分离的关键。
Engine 提供 FlushSignal 能力,但不会决定所有 Sink 收到 Signal 后必须执行什么动作。
因为不同 Sink 的 flush 语义并不相同。
对于 JDBC Sink,flush 可能意味着执行批量 SQL 写入;对于其他 Sink,flush 可能对应不同的数据提交逻辑。因此,Connector 需要根据自身实现决定是否注册 flushAction。
这种设计让 FlushSignal 成为一种可扩展的 Engine 能力,而不是绑定某一个 Connector 的特殊机制。
最终,Flush 从一个 JDBC Sink 内部参数,演变成 SeaTunnel Engine 中统一管理的控制事件。Engine 负责 Signal 的产生和传播,Connector 负责 Signal 到达后的具体执行,两者通过明确的边界完成协作。
这也是 STIP-23 对 SeaTunnel Zeta Engine 的一次重要抽象提升。它解决的不只是 batch_interval_ms 的实现问题,而是让 Engine 具备了更加完善的运行时控制能力。
05 Exactly-Once 约束
前面通过 FlushSignal 将定时 Flush 能力从 Connector 层提升到了 Engine 层,但新的问题也随之出现:当 FlushSignal 参与数据处理流程后,如何保证原有的 Exactly-Once 语义不被破坏。
对于支持 XA 事务的 Sink 来说,Flush 并不等价于事务提交。FlushSignal 可以触发数据刷新或者事务准备,但不能绕过 Checkpoint 机制直接决定事务最终提交。因为在 SeaTunnel 的 Exactly-Once 模型中,事务状态必须与 Checkpoint State 保持一致,只有进入 Checkpoint 管理范围内的事务,才具备最终提交资格。
因此,Engine-Level FlushSignal 的设计必须满足一个核心约束:FlushSignal 可以推动事务状态变化,但不能脱离 Checkpoint 独立完成 Commit。
在引入 FlushSignal 之前,Sink 的事务提交流程由 Checkpoint 控制。Checkpoint Barrier 到达后,系统保存当前事务状态,并在 notifyCheckpointComplete阶段完成最终提交。
如果 FlushSignal 直接触发 XA COMMIT,就会打破这一流程。

图中左侧的方式不推荐,因为这种方式的问题在于,事务提交发生在 Checkpoint State 记录之前。
例如,一个事务对应的 XID 已经完成 XA COMMIT,但是此时 Checkpoint 尚未完成保存,那么从 SeaTunnel 状态管理角度来看,这个事务并不存在于可恢复状态中。
如果此时 Source 发生故障并从之前的 Checkpoint 恢复,Source 可能重新消费已经提交的数据,从而产生重复写入。
因此,Timer 或 FlushSignal 都不能直接承担事务提交职责。
正确的设计方式如上图右侧所示。在这种模式下,FlushSignal 只负责推动当前事务进入 prepare 状态,而真正的 Commit 仍然等待 Checkpoint 完成。
也就是说,PREPARE 可以提前发生,但 COMMIT 必须等待 CK。
这样,FlushSignal 虽然改变了事务处理节奏,但不会改变 Exactly-Once 的控制边界。
基于这一原则,当前 JDBC XA Writer 实现中并不会注册 Timer Flush,而是继续让事务提交受到 Checkpoint 生命周期管理。
在保证 FlushSignal 不直接提交事务之后,另一个问题是:FlushSignal 如何影响 XA 事务的划分。
当前设计方案(非当前能力):

FlushSignal 到来后,并不会结束整个 Checkpoint 生命周期,而是帮助 Sink 提前切分事务范围。
根据设计,事务可以划分为:
txn-1R R R F1txn-2R R F2txn-3R R CK
当 F1 到达时,Sink 可以执行:
executeBatch()XA PREPAREbeginTx()
完成当前阶段的数据准备。
随后新的数据进入新的事务范围,直到下一次 FlushSignal 到来。
F2 到达后,会继续触发下一次 prepare。
但是这些已经 prepare 的事务,并不会立即提交。
pendingCommitInfo / pendingStates 先留在内存,CK 时统一合并并持久化。
也就是说,FlushSignal 只负责帮助 Sink 切分事务边界,将较大的事务拆分成多个 prepared transaction,但这些事务仍然等待 Checkpoint 统一确认。
当 Checkpoint 到达时,系统会执行:
prepare txn-3snapshotState()merge pending↓Checkpoint State[txn-1, txn-2, txn-3]
Checkpoint State 保存了当前所有已经 prepare 的事务信息。
随后,在:
notifyCheckpointComplete
阶段执行:
XA COMMIT ALL
完成最终提交。
这种设计使 FlushSignal 和 Exactly-Once 可以同时存在。FlushSignal 提供更灵活的事务切分能力,而 Checkpoint 仍然掌握最终提交权。
在分布式运行环境中,任务可能在任意阶段发生故障。因此,除了正常提交流程之外,还需要明确恢复时应该依据什么判断事务状态。
核心判断标准是:恢复依据始终是 XID 是否已经进入可恢复的 Checkpoint State。
FlushSignal 本身并不是恢复依据,它只是影响事务 prepare 的时间点。
根据故障发生的位置,可以分为三种恢复情况。

第一种情况是 Checkpoint 到达之前发生故障。
此时虽然 F1、F2 对应事务已经 prepare,但由于它们还没有进入 Checkpoint State,因此恢复时无法认为这些事务已经属于系统状态。
恢复过程:
State 无 txn-A / txn-B↓XA RECOVER↓ROLLBACK↓Source 从 CK-N 重放
这些事务会被回滚,然后由 Source 根据旧 Checkpoint 重新消费数据。
第二种情况是事务状态已经进入 Checkpoint State,但 Commit 前发生故障。
此时 XID 已经保存到 Checkpoint State。
恢复时:
State 含 txn-A / txn-B / txn-C↓restoreCommit↓COMMIT
系统可以根据 Checkpoint 中保存的信息继续完成事务提交。
第三种情况是 Checkpoint 已经完成提交后发生故障。
恢复时:
事务已提交↓XA RECOVER 扫描为空↓no-op
系统无需再次执行操作。
通过这三种恢复场景可以看到,FlushSignal 并没有改变 SeaTunnel 原有的 Exactly-Once 机制,而是在 Checkpoint 管理框架内增加了一种更灵活的事务触发方式。
FlushSignal 负责触发 prepare,Checkpoint State 负责记录事务状态,notifyCheckpointComplete 负责最终提交。
这种职责划分保证了 Engine-Level Flush 能力不会破坏 XA 事务一致性,同时也为未来更多运行时控制事件提供了统一设计基础。
06 机制与职责边界
从 JDBC batch_interval_ms 的问题演进到 Engine-Level FlushSignal,实际上也是一次对 SeaTunnel Engine 和 Connector 职责边界的重新认识。

最开始,如果只从 JDBC Connector 的角度思考,解决方式很容易落在 Connector 内部。
因为问题表面上是 JDBC Sink 需要定时 flush。于是我们自然会想到 JDBC 自己创建 Timer。
但是进一步分析之后可以发现,Timer 管理、线程协调、异常传播、任务生命周期处理,这些内容并不是 JDBC 特有的问题。它们属于数据处理引擎运行时需要提供的基础能力。
如果每一个 Connector 都自己实现一套 Timer 机制,那么随着 Connector 数量增加,系统会出现大量重复逻辑。
每个 Connector 都需要考虑自己的后台线程如何启动,如何停止,如何处理异常,如何避免与 checkpoint、close、cancel 等流程产生冲突。
最终,一个简单的数据写入参数,会迫使 Connector 开始维护越来越复杂的运行时逻辑。
这显然不是合理的职责划分。
在 FlushSignal 的设计中,SeaTunnel 重新明确了 Engine 和 Connector 各自应该负责的内容。
Connector 关注的是外部系统本身的语义。
例如 JDBC Sink 需要知道如何执行 batch 写入,SQL 如何执行,transaction 如何 prepare、commit 或 rollback,以及收到 flush 请求之后具体应该做什么。
这些逻辑只有 Connector 自己最清楚。
而 Engine 负责的是更加通用的运行机制。
例如 Timer 生命周期如何管理,如何产生控制事件,如何让 Signal 进入数据流,以及如何保证 Signal 最终在正确的 Task 执行路径中被消费。
这些能力并不属于某一个具体 Connector,而属于整个数据处理引擎。
因此,更合理的设计不是让 JDBC 自己维护 Timer,而是:
Connector 注册 flushAction。
Engine 负责产生 FlushSignal,并将 Signal 传递到 Sink。
最终由 Sink 在自己的消费线程中执行 flush。
这样,时间触发机制和具体业务行为被进行了分离。
在 STIP-23 的设计中,Engine 提供了相应的配置能力:
sink.flush.interval = 5000
这个配置表达的是 Engine 层的定时行为。
当配置开启后,Engine 会按照指定时间间隔尝试产生 FlushSignal。
默认情况下,如果配置为 0,则表示关闭该能力,不会影响已有任务。
但是 Engine 并不会直接调用某个 Connector 的 flush 方法。
因为 Engine 并不知道 JDBC 的 flush 是 executeBatch(),也不知道其他 Sink 的 flush 是 bulk write、stream load,还是其他内部操作。
因此,Connector 需要通过 Context 注册自己的 flushAction。
例如:
context.registerFlushAction(() -> { flush();});是否注册,由 Connector 自己决定。
在整个设计中,最关键的一点是 Timer 不直接调用 Sink。
如果 Timer 线程直接执行:
Timer Thread → SinkWriter.flush()
那么实际上又回到了最初 Connector 后台线程的方案。之前存在的问题依然存在:
因此,STIP-23 选择让 Timer 只负责产生 FlushSignal。Signal 进入 SeaTunnel 正常的数据通道之后,最终由 SinkFlowLifeCycle 在消费线程中处理。
整个过程变成:
Timer↓FlushSignal↓Data Path↓SinkFlowLifeCycle↓flushAction.run()
这样 flush 又回到了 SeaTunnel 原本的数据处理模型中。
07 思考与总结:从 FlushSignal
看 SeaTunnel Engine 的演进方式
从 JDBC batch_interval_ms 到 Engine-Level FlushSignal,这个过程表面上是在解决一个定时 flush 问题,实际上反映的是数据引擎在能力扩展时如何进行抽象设计。一个功能能否长期演进,并不取决于实现代码多少,而取决于是否明确语义边界、是否复用已有架构,以及是否提供合理的扩展机制。
在设计 FlushSignal 时,首先需要明确 flush 的语义。FlushSignal 并不代表一次提交成功,也不代表数据已经对外可见,它代表的是 Engine 向 Sink 提供了一次执行 flush 的机会。
不同系统对于“成功”的定义并不相同。At-least-once 允许失败恢复时重复处理,而 Exactly-once 则需要依靠 Checkpoint、事务等机制保证结果只产生一次。因此,FlushSignal 只能负责触发行为,而不能替代 Connector 自身的一致性控制。
这也是为什么 Engine 不应该强制所有 Sink 执行 flush。对于不同 Connector 来说,flush 可能意味着 batch 写入、事务准备,或者其他特殊操作。如果 Engine 不区分这些语义,反而可能破坏原有的事务和 Exactly-once 保证。
因此,一个通用能力首先需要定义清楚“提供什么能力”和“不负责什么问题”。

FlushSignal 的设计没有引入新的执行链路,而是在现有 Task 架构中增加了一种事件类型。
如果 Timer 直接调用 Sink:
Timer Thread → SinkWriter.flush()
虽然实现简单,但会产生新的线程模型,导致 flush 与数据写入、checkpoint、close 等生命周期操作之间出现并发问题,同时后台线程异常也无法自然回传到 Task 执行流程。
因此,FlushSignal 选择复用已有的数据通道。Engine Timer 触发后,将 Signal 注入 Source 输出,随后沿着 Source → Transform → Sink 的链路传递。
在这个过程中,Source 负责注入 Signal,Transform 只负责透传,SinkFlowLifeCycle 最终识别 Signal 并执行对应动作。这样,FlushSignal 和 Record 共享同一套运行模型,不需要额外维护新的线程和控制路径。
这种设计体现了 Engine 演进的重要原则:新增能力应该尽量融入已有抽象,而不是不断增加特殊逻辑。

FlushSignal 的价值不仅在于解决 JDBC flush 问题,更重要的是体现了 Engine 和 Connector 的职责划分。
Engine 负责通用运行机制,包括 Timer 管理、Signal 生成、生命周期协调以及数据流传递;Connector 负责具体业务语义,包括 flush 如何执行、是否支持定时 flush,以及如何保证自身事务一致性。
通过 Context 或 SPI 方式注册 flushAction,Connector 可以选择是否使用该能力,Engine 不需要了解具体实现细节。
这种设计带来了三个优势:默认行为不会影响已有 Connector,新增能力不会破坏原有语义,同时未来类似的运行时控制需求也可以基于相同模型扩展。
回到最初的 batch_interval_ms 问题,它并不是 JDBC 一个参数实现的问题,而是暴露了 Engine 缺少运行时控制事件的问题。从 Connector 内部线程,到 Engine-Level Signal,这实际上是一次从局部优化走向架构抽象的过程。
一个成熟的数据引擎,并不是把所有功能都提前实现,而是通过合理的抽象和扩展点,让新的能力能够以更低成本、更稳定的方式融入系统。

Apache SeaTunnel是一个云原生的多模态、高性能海量数据集成工具。北京时间 2023 年 6 月1 日,全球最大的开源软件基金会ApacheSoftware Foundation正式宣布SeaTunnel毕业成为Apache顶级项目。目前,SeaTunnel在GitHub上Star数量已达9k+。SeaTunnel支持在云数据库、本地数据源、SaaS、大模型等170多种数据源之间进行数据实时和批量同步,支持CDC、DDL变更、整库同步等功能,更是可以和大模型打通,让大模型链接企业内部的数据。
同步Demo
新手入门

最佳实践

测试报告

源码解析



