跳到主要内容

4 篇博文 含有标签「SeaTunnel」

查看所有标签

· 阅读需 10 分钟
David Zollo

问一位数据工程师他的管道是 ETL 还是 ELT,几乎总能得到秒答——老派工程师说 ETL,dbt 用户说 ELT。

但这两个答案都不完整。有第三种模式更准确地描述了现代数据管道实际在做的事——它一直存在,只是没有被命名:EtLT


三种范式

ETL:落库前先转换

原始数据从源端提取后,经过专用转换层(Spark、DataStage、Informatica)处理,再写入目标数仓。

源端 ──► [转换层] ──► 目标端

优点: 目标端始终存放干净、业务就绪的数据。合规管控在数据落库前就已强制执行。
缺点: 转换层容易成为瓶颈。Schema 变更需要跨多层同步修改。专用计算集群的运维成本高。

ELT:先落库,原地转换

由 dbt 带火的模式。原始数据直接落入数仓(BigQuery、Snowflake、ClickHouse),再用 SQL 原地做转换。转换层不再是独立集群——就是数仓本身。

源端 ──► 目标端(原始层)──► [数仓内 SQL] ──► 业务层

优点: 原始数据保留完整审计价值。复用数仓算力。迭代速度快。
缺点: 手机号、身份证、邮箱等敏感字段以明文落库,存在合规风险。数据质量问题要等 Load 之后才可见,此时下游模型可能已经被污染。

EtLT:传输途中做轻量变换

源端 ──► [tiny t] ──► 目标端(净化原始层)──► [目标端内的 T]

tiny t 是数据在传输过程中完成的一组行级变换:

tiny t 操作目的
字段投影 / 列裁剪传输前删除无用列——节省带宽
PII 脱敏手机号、身份证、邮箱在落库前匿名化——合规在管道层强制落地,而非事后补救
类型归一源端 VARCHAR "2023-01-01" 在入仓时就转为 DATE——省去数仓侧的类型转换 SQL
行过滤不需要的事件可在 Load 前丢弃;有状态的 CDC 变化需要按 changelog 语义处理
字段重命名适配目标端命名规范,无需数仓侧的别名层
NULL 回填减少下游聚合 SQL 中的 COALESCE 调用

大 T(Load 后的 Transform)才是真正的业务逻辑所在:多表 JOIN、指标计算、聚合、ML 特征工程。


为什么工程师长期忽略 EtLT

根本原因是工具,不是概念本身。

传统 ETL 工具(Informatica、DataStage)把转换做得极重,工程师被迫把大量业务逻辑塞进管道,结果管道变得脆弱、难以维护。

现代 ELT 工具(Airbyte、Fivetran)走向另一个极端:源端到目标端几乎没有在途处理,全部交给 dbt。

两类工具都没能填好的空白: 当你同时需要在途操作(脱敏、过滤)目标端的复杂分析 SQL 时,两个工具族都会让你陷入尴尬的临时方案。

EtLT 天然填补这个空白。


SeaTunnel 天然契合 EtLT

Apache SeaTunnel 在官方文档中把自己定位为 "EL(T) 数据集成工具"——括号里的 T 是刻意为之。Transform 步骤是轻量的、可选的。这不是营销话术,而是一个架构约束。

三层模型与 EtLT 直接对应

SeaTunnel 的执行模型分三层:Source → Transform → Sink

Source(E)
└──► Transform(tiny t) ← 可选,轻量处理
└──► Sink(L)
└──► [数仓 SQL / dbt](T) ← 大 T 在这里

SeaTunnel 对 Transform 能做的事划定了清晰边界。官方文档写道:

"Transform can only be used for some simple transformations of data, such as converting a column to uppercase/lowercase, modifying column names, or splitting one column into multiple columns."

这描述的是内置 Transform 的目标边界。当前 Zeta SQL Transform 不支持 JOIN 和 GROUP BY,因此跨表 JOIN 与聚合应放在目标系统或其他专用处理层中完成。

内置 Transform 覆盖典型 tiny t 需求

SeaTunnel Transform对应的 tiny t 操作
FieldMapper字段重命名、列投影
Filter通过 include/exclude 列表做字段投影
Replace字段值替换(脱敏 / 打码)
Split将一列拆成多列(如地址解析)
SQL Transform行过滤和轻量 SQL 表达式;当前 Zeta 实现不支持 JOIN 或 GROUP BY
Copy字段复制

SeaTunnel 同时提供 map 和 flat-map Transform 接口。当前内置 Transform 主要处理单条记录和 Schema,而不是跨行聚合,因此很适合作为 tiny t 层。

EtLT 最关键的战场:CDC 管道

实时 CDC 同步是 SeaTunnel 的核心使用场景之一——也是纯 ELT 在合规上失守的地方。

根本问题在于:binlog 事件里的字段值一旦落库,就没有第二次机会打码了。 写入 ClickHouse 的手机号或身份证,不可能被下游 dbt 模型"反写"回来。原始明文已经持久化。

EtLT 解决了 ELT 无法解决的问题:在数据到达目标端之前的唯一时间窗口内,完成合规变换。

下面给出一个带 tiny t 的 SeaTunnel CDC 作业最小配置形态。运行前请替换示例中的账号、密码和服务地址:

env {
job.mode = "STREAMING"
}

source {
MySQL-CDC {
plugin_output = "raw_user_info"
url = "jdbc:mysql://localhost:3306/orders"
username = "seatunnel"
password = "change-me"
server-id = 5601-5604
table-names = ["orders.user_info"]
}
}

transform {
# 在传输途中打码手机号——原始值永远不会到达数仓
Replace {
plugin_input = "raw_user_info"
plugin_output = "masked_user_info"
replace_field = "phone"
pattern = "(\\d{3})\\d{4}(\\d{4})"
replacement = "$1****$2"
is_regex = true
}
}

sink {
# L:落入 ClickHouse 原始层
Clickhouse {
plugin_input = "masked_user_info"
host = "clickhouse-host:8123"
database = "raw"
table = "user_info"
username = "default"
password = "change-me"
primary_key = "id"
support_upsert = true
allow_experimental_lightweight_delete = true
}
}

数据进入 ClickHouse 后,再用 SQL 构建宽表、计算留存指标、做分析查询——这就是大 T。dbt 模型在这里是天然的搭档。


SeaTunnel + dbt:全栈 EtLT

dbt 负责 Load 之后的大 T;SeaTunnel 负责从源到 Load(含 tiny t)。两者互补,而非竞争。

源端
└── SeaTunnel(E + tiny t + L)
└── dbt(T)
└── BI / ML

SeaTunnel 负责数据到达和在途卫生;dbt 负责数据建模和业务逻辑。每个工具只做它擅长的事。


常见陷阱:把大 T 塞进 SeaTunnel

因为 SeaTunnel 支持 SQL Transform,有些工程师会尝试在里面写 GROUP BY 聚合。当前 Zeta SQL Transform 会明确拒绝 GROUP BY 和 JOIN 查询。

如果确实需要 Load 前的聚合,有两条路:

  1. 在目标系统里做——这本来就是 EtLT 和 ELT 中大 T 的用途。
  2. 在 SeaTunnel Transform 链路之外使用专用处理作业——此时构建的是完整流式 ETL 管道,而不是把转换留在 tiny t 层。

边界清晰,失败才容易定位。把 tiny t 和大 T 混在同一层,出问题时几乎无法推断根源在哪里。


总结

范式最适合的场景SeaTunnel 的角色
ETL强合规要求、有专用转换计算集群负责 E、L 和已支持的轻量转换;复杂 T 需要专用处理层
ELT目标端是强 SQL 引擎(BigQuery/Snowflake),无严格在途合规要求纯 E + L——关闭 Transform
EtLT实时 CDC、在途 PII 脱敏、目标端有 dbt/SQL 能力天然契合:E + tiny t + L;大 T 归目标端

SeaTunnel 天然适合 EtLT——不是因为它声称这个标签,而是因为 Source/Transform/Sink 三层分离、当前内置 Transform 的能力边界和明确的"EL(T)"定位都指向同一件事:快速移动数据,在传输途中轻量净化,把重活留给目标系统。


延伸阅读

· 阅读需 13 分钟
Daniel

同程旅行的数据通道经过多年演进,逐步形成了离线搬运、实时集成、Sqoop 和 SeaTunnel 四套系统并存的格局。每套系统都解决了特定阶段的问题,但功能重叠、执行引擎割裂、运维体系分散,也逐渐成为平台统一治理的障碍。

在 Apache SeaTunnel Meetup 上,负责同程旅行数据平台相关工作的周晓晨,分享了同程旅行以 Apache SeaTunnel Zeta 引擎为统一底座,将四套系统收敛为批流一体数据通道的实践。整个项目有三个不可妥协的目标:迁移过程对业务透明、切换前证明数据一致、同时提升执行效率和运行稳定性。

本文以该场公开分享的事实为依据,独立梳理其统一架构、迁移保障、AI 辅助任务生成、数据校验设计以及后续规划。

· 阅读需 19 分钟
David Zollo

过去二十年,企业数据工程体系一直建立在一个默认前提之上:

人负责理解系统,系统负责执行流程。

工程师理解业务,拆解链路,编写 SQL、Spark、Shell、同步脚本或调度任务,然后交给系统执行。调度系统不需要理解业务,只需要按照 DAG 把任务跑完;数据同步系统不需要理解指标,只需要把数据从源端搬到目标端。

这套模式支撑了很长时间的数据仓库、数据湖、BI、报表和批量调度体系。

但今天,这个前提正在开始失效。

企业数据系统变得越来越复杂:数据源更多,链路更长,实时性更强,业务变化更快,指标口径更容易冲突,AI 应用又不断带来新的交互数据、模型反馈数据、向量索引和非结构化数据。

在这样的环境下,企业真正缺少的,已经不只是更多 Pipeline,也不只是一个更会写 SQL 的 Copilot。

企业真正需要的是:

能够理解系统、规划任务、调用工具、验证结果,并持续积累经验的 Data Engineering Agent。

而在这个变化中,Apache SeaTunnel 的角色会变得非常关键。

因为 Agent 时代的数据工程,不只是“会思考”就够了。它还必须能稳定连接数据源、捕获变化、执行同步、处理增量、保障一致性,并把数据安全、可靠、低成本地移动到目标系统。

换句话说:

Agent 负责理解目标和规划动作,SeaTunnel 负责把这些动作变成真实、可靠、可恢复的数据工程执行。

这也是为什么从 ETL、ELT、EtLT 到 Agent 的演进里,SeaTunnel 会成为下一代数据工程体系的重要执行底座。

一、从 ETL 到 ELT:第一次范式迁移

传统 ETL 的逻辑很清楚:

  1. 先 Extract,把数据从源端抽出来。
  2. 再 Transform,在中间层完成清洗、转换、聚合。
  3. 最后 Load,把处理后的结果加载到目标系统。

这套模式适合早期数据仓库时代。

因为当时数据源相对有限,链路相对清楚,计算资源也比较集中。企业希望在进入数仓之前,把脏数据清洗掉,把口径处理好,把结构整理成可以直接使用的结果。

ETL 本质上是一种确定性 Pipeline。

它的核心假设是:

人提前定义好流程,系统按流程执行。

后来,随着云数仓、数据湖、Lakehouse 和弹性计算的发展,ELT 开始流行。

ELT 的思路是:

  1. 先 Extract。
  2. 再 Load。
  3. 最后在目标系统里 Transform。

也就是说,先把原始数据或近原始数据装进统一存储底座,再利用目标端强大的计算能力做转换。

ELT 解决了 ETL 的一些问题。比如,它降低了前置处理的复杂度,保留了更多原始数据,也让后续建模和分析更加灵活。

但 ELT 也带来了新的问题。

如果所有转换都推迟到目标端,源端数据的脏问题、类型问题、字段变化、CDC 事件、隐私字段、格式差异、结构漂移,都会被直接带进目标系统。

在简单批量场景里,这还可以接受;但在今天的实时同步、CDC、多表同步、湖仓入湖、SaaS API 接入和 AI 数据工程场景里,完全“先 Load 再说”会让下游治理成本急剧上升。

于是,企业数据工程开始进入第三种形态:

EtLT。

二、EtLT:从 Pipeline 到 Agent 的关键过渡形态

这里的 EtLT,不是简单把 ETL 和 ELT 混在一起。

我更愿意把它理解为:

Extract -> lightweight transform -> Load -> semantic Transform

也就是:

  • 抽取数据。
  • 先做一层必要的轻量 transform。
  • 再把数据加载到统一数据底座。
  • 最后做面向业务语义、指标口径和数据产品的大 Transform。

这里最关键的是小写 t 和大写 T 的区别。

小写 t 不是复杂业务建模,而是数据进入平台之前必须完成的工程化处理,例如:

  • 字段裁剪
  • 类型映射
  • 格式标准化
  • 主键或分区字段处理
  • 敏感字段脱敏
  • CDC 事件格式转换
  • 多表路由
  • Schema Evolution
  • 数据质量前置校验
  • 一读多写
  • 限速和并行控制

这些处理不应该全部推迟到目标端。否则,数据湖、数仓或 Lakehouse 里会堆积大量结构不一致、语义不清楚、质量不可控的原始数据。

但小写 t 也不应该承载过重的业务逻辑。

真正复杂的业务口径、指标定义、主题建模、语义映射和跨域聚合,应该放到大写 T 里,在 Lakehouse、数仓、指标层、语义层和治理体系中完成。

所以 EtLT 的核心价值在于:

在数据进入统一底座之前,先完成必要的工程标准化;在数据进入统一底座之后,再完成业务语义化。

这正是 SeaTunnel 非常适合承担的部分。

Apache SeaTunnel 的 Source、Transform、Sink 架构,本质上天然适合 EtLT 里的小写 t:它既能连接多种异构数据源,也能在数据流动过程中完成轻量转换、结构适配、CDC 处理、多表同步和写入目标系统。

因此,在 EtLT 架构里,SeaTunnel 不只是一个“数据搬运工具”,而是企业数据进入统一底座之前的 Data Integration Runtime

三、为什么传统 ETL 开始吃力

传统 ETL 面对的是相对确定的流程。

把规则写清楚,DAG 画出来,任务按顺序执行,失败之后工程师排查修复。

但今天的数据工程环境已经不是当年的样子。

企业里同时存在 OLTP 数据库、Kafka 消息流、CDC 链路、SaaS API、日志、对象存储、Lakehouse、实时 OLAP、向量数据库、AI 交互日志和模型结果数据。

数据不只是变多了,也变得更碎、更异构、更实时。

更麻烦的是链路长度。

一个看上去很普通的经营指标,背后可能经过几十张表、多层宽表、多个业务域拼接,还有一堆复杂口径转换。到了这个阶段,很多企业面临的真实问题已经不是“流程有没有搭起来”,而是几乎没人能把整条链路从头到尾讲明白。

传统 ETL 在这里暴露出来的是结构性问题。

  • 一个字段改名,可能带崩上百个任务。
  • 一个枚举值变化,可能让多个核心指标悄悄漂移。
  • 一段增量逻辑写错,影响的可能不只是单张表,而是一串下游分析结果。

调度系统能告诉你任务失败了,但它不一定知道为什么失败。 同步工具能把数据写过去,但它不一定知道这份数据影响了哪个业务指标。 工程师能修复脚本,但前提是他先能找到完整上下文。

所以传统 ETL 真正吃力的地方,不只是性能或稳定性,而是它天然缺少系统级认知能力。

它能执行流程,但不能理解系统。

四、为什么 Copilot 没有真正解决问题

很多团队第一次把 AI 引入数据工程,往往是从 Copilot 开始:写 SQL、补 Spark 代码、生成 YAML、写测试样例。

这些能力当然有价值,尤其能提高局部开发效率。

但它们没有触达企业数据工程最深层的问题。

因为数据工程真正难的,从来不只是代码生成,而是系统理解。

Copilot 可以帮你生成一段 SQL,却不知道这个字段的业务含义是什么;可以帮你写一个同步任务,却不知道上游 schema change 会影响哪些下游指标;可以帮你生成调度配置,却不知道这次变更是否破坏了历史一致性。

企业数据工程真正复杂的部分,是:

  • lineage reasoning
  • dependency analysis
  • semantic understanding
  • metric governance
  • risk estimation
  • impact analysis
  • incremental recovery

这些不是单靠代码补全就能解决的。

所以企业真正需要的,不只是一个 AI IDE,而是一套能理解数据系统、围绕目标拆解任务、调用工程工具并验证结果的 Agentic Data Engineering 系统。

五、从 Pipeline 到 Agent,真正变化是什么

如果只保留一句核心判断,我会这样说:

传统 ETL 的核心是:人定义流程,系统执行流程。Agent 数据工程的核心是:人定义目标,系统生成流程。

这不是一句包装口号,而是两套系统组织方式的根本差异。

在传统模式里,工程师先设计好任务链路,配置 Source、Transform、Sink,再交给调度系统执行。

系统面对的是一个固定流程。

而在 Agent 模式里,人输入的可能只是一个目标。

比如:

新增一个订单利润指标,并保证和财务口径一致。

传统做法是工程师自己去找数据源、查表结构、看血缘、写转换逻辑、配置同步任务、补质量校验、发起回归测试。

Agent 理想中的工作方式,则是围绕这个目标自动拆解一组动作:

  • 识别涉及哪些业务对象
  • 查找候选数据源
  • 分析上游 lineage
  • 判断应该走 ETL、ELT 还是 EtLT
  • 生成 SeaTunnel 同步任务
  • 配置 CDC 或全量同步
  • 完成轻量 transform
  • 写入 Lakehouse 或目标系统
  • 触发质量校验
  • 评估下游影响
  • 把结果反馈给工程师确认

这里真正变化的,不是“AI 帮你写了一段 SQL”,而是系统开始围绕目标生成工程动作。

但这也带来一个关键问题:

Agent 规划出来的数据动作,谁来可靠执行?

这正是 SeaTunnel 的价值所在。

六、SeaTunnel 在 Agent 时代的位置:Agent 的数据集成执行层

Agent 不能只停留在分析和建议阶段。

如果一个 Data Engineering Agent 发现某张表需要同步、某条链路需要回放、某个 CDC 任务需要调整、某批数据需要重新写入目标系统,它必须能调用一个稳定、可控、可观测的数据集成执行层。

这个执行层需要具备几类能力。

1. 能连接足够多的数据源

企业数据系统天然异构。Agent 不可能只面对一种数据库或一种文件系统。它需要连接 MySQL、Oracle、PostgreSQL、SQL Server、Kafka、Hive、Iceberg、Doris、ClickHouse、StarRocks、Elasticsearch、S3、HDFS、MongoDB、SaaS API 等不同系统。

SeaTunnel 的 Connector API 和插件化架构,正好解决了这个问题。它把 Source、Transform、Sink 抽象成统一接口,让不同数据源以相对一致的方式被接入和调用。

2. 能同时处理批、流、CDC 和整库同步

Agent 时代的数据工程不是单一范式。它既需要一次性全量迁移,也需要持续 CDC;既需要离线批处理,也需要实时同步;既需要单表同步,也需要多表或整库级别的数据移动。

SeaTunnel 的价值就在于,它不是单纯的 ETL 脚本工具,而是面向数据同步和数据集成场景构建的运行体系,可以承载全量、增量、实时、CDC、多表同步等多种数据移动任务。

3. 能承载 EtLT 里的轻量 transform

EtLT 不是把所有业务逻辑都塞进同步链路,而是在数据进入统一底座之前,先完成必要的工程标准化。

SeaTunnel 的 Transform 能力,适合处理字段映射、类型转换、数据过滤、字段裁剪、格式调整、脱敏、路由等轻量处理。

这些能力让 SeaTunnel 可以承担 EtLT 中的小写 t

不做过重的业务建模,但把数据整理到可治理、可加载、可继续加工的状态。

4. 能提供一致性、容错和恢复能力

Agent 可以判断“应该重跑这一段链路”,但真正执行重跑时,底层系统必须具备 Checkpoint、失败恢复、状态管理和一致性保障能力。

一个只会规划、不可靠执行的系统,最终只是“会想但不会做”。

所以在 Agent 时代,执行质量仍然和智能水平同样重要。

七、未来的数据工程栈会越来越像“操作系统”

如果再往前看一步,企业数据工程体系会越来越像一套分层的数据操作系统,而不再只是很多零散 Pipeline 的集合。

在这套体系里:

  • Semantic Layer 负责定义业务世界模型
  • Metadata 负责提供结构化上下文
  • Memory 负责沉淀经验
  • Planning Layer 负责把目标拆成动作
  • Execution Layer 负责真正执行同步、CDC、数据移动、回放和恢复

SeaTunnel 就处在这个 Execution Layer 里。

这个位置非常关键。

未来不是“在 ETL 上面套一个大模型”就结束了。

未来是一个清晰分层、各司其职的系统:

  • Agent 决定应该做什么
  • SeaTunnel 保障这件事真的被可靠执行出来

八、用一句话总结这个演进

ETL 时代,企业构建的是数据流水线。

ELT 时代,企业开始把数据沉淀到统一底座。

EtLT 时代,企业开始重新平衡数据接入、轻量治理和业务转换之间的关系。

Agent 时代,企业要构建的更像是一套数据操作系统。

在这套系统里,Semantic Layer 提供业务世界模型,Metadata 提供上下文,Memory 沉淀经验,Planning Layer 负责决策,SeaTunnel 这样的数据集成引擎负责把数据移动、CDC、轻量 transform、增量同步、回放和恢复真正执行出来。

所以,Agent 不是在 ETL 外面简单挂一个大模型,也不是给调度系统加一层聊天入口。

它真正改变的是企业数据工程的底层操作范式。

而 SeaTunnel 的价值,也会在这个过程中被重新定义。

它不只是一个数据同步工具,而是下一代 Agentic Data Engineering 体系里的数据集成执行底座。

Agent 让数据系统开始理解目标。

EtLT 让数据进入平台的过程更加可控。

SeaTunnel 让这些目标真正变成可靠的数据工程动作。

这就是从 ETL、ELT、EtLT 到 Agent 时代,企业数据工程正在发生的根本变化。

· 阅读需 7 分钟

在数据集成与同步领域,Apache SeaTunnel 无疑是当下最炙手可热的工具之一。本系列将深挖其高级用法。

首篇从 SeaTunnel 核心概念“数据流”切入,剖析底层原理,如数据流动与转换机制,结合实例讲解在复杂场景中的应用,助你掌握这一工具。

一句话总结(先给结论)

SeaTunnel 不是“source → sink”的线性工具

  • 它是一个 “数据流(DataStream / DataFlow)驱动的 DAG 执行引擎”

两个 source 可以流入一个 sink,正是这个模型的直接体现。

一、SeaTunnel 的核心概念:数据流

在 SeaTunnel 内部,一切围绕“数据流”展开

数据流是什么?

数据流 = 一组结构一致的 Record 流(带 Schema)

它不是表、不是文件、不是 SQL 结果 而是:

Record1 → Record2 → Record3 → ...

每个插件都在“操作数据流”

二、plugin_output / plugin_input 的真实含义(非常重要)

你之前一直在“用”,但现在该“理解”它了。

1️⃣ plugin_output

plugin_output = "source_data_output_1"

含义不是“名字”,而是:

给当前插件产生的数据流起一个唯一 ID

可以理解为:

DataStream<ID = source_data_output_1>

2️⃣ plugin_input

plugin_input = "source_data_output_1"

含义是:

我这个插件,要消费哪个数据流

用一句话说透

plugin_output / plugin_input = 数据流的“连线端口”

三、SeaTunnel 的 DAG 模型(你现在已经用到了)

你这个成功的实验,本质上是:

SourceA ─┐
├──► Sink
SourceB ─┘

SeaTunnel 内部会构建这样的 DAG:

DataStream A ─┐
├──► Sink Operator
DataStream B ─┘

关键点:为什么能合并?

因为:

Sink 不是“绑定一个 source”,而是“订阅一个或多个数据流”

当你写:

sink {
jdbc {
plugin_input = "a,b"
}
}

或你用了多个 source 但最终都连到同一个 sink,SeaTunnel 内部会:

  • 把多个输入流
  • 合并成一个 Logical Input
  • 按 Record 顺序写入

⚠️ 注意:

  • 不是 join
  • 不是 union SQL
  • 是流级别的合并(append)

四、这和你理解的 “SQL / ETL” 有什么本质区别?

这是很多人会迷糊的地方。

SQL 的世界

SELECT * FROM A
UNION ALL
SELECT * FROM B
  • 是 “结果集语义”

SeaTunnel 的世界

A 的 Record 流
B 的 Record 流

Sink 持续消费
  • 是 “流语义”

只要 Schema 一致,就能进同一个 sink。

五、Schema 在数据流里的地位(你一定要记住)

数据流 = Record + Schema

SeaTunnel 合流的前提:

  • 字段数一致
  • 字段类型兼容
  • 字段名对齐(或可映射)

否则:

  • 运行期直接报错
  • 或 Sink 写入失败

你刚才说“目标字段肯定是对齐的”,这就是为什么你的实验能成功。

六、SeaTunnel 的“数据流转模型”正式定义(给你一个标准版)

你以后在设计 / 讲方案 / 写文档,可以直接用这段话:

SeaTunnel 以 数据流(DataStream) 为核心抽象。 Source 插件负责生成数据流,Transform 插件对数据流进行处理并输出新的数据流,Sink 插件消费一个或多个数据流并将数据写入外部系统。 多个数据流可以在 Sink 处汇聚,只要其 Schema 兼容,SeaTunnel 会以流合并(append)的方式进行写入,而不是关系型 join。

七、这对你 Builder / Strategy 设计的直接影响(重点)

你现在可以非常确定三件事:

1️⃣ Builder 必须支持 N Source → M Sink

不是 1→1,而是 图模型

2️⃣ plugin_output 是“一等公民”

你现在 Builder 里如果有人不设置 plugin_output

  • 你应该帮他自动生成一个

这是平台级能力。

3️⃣ Sink 允许多个输入流(逻辑)

即使 DSL 上写的是一个:

plugin_input = "s1"

你 Builder 的语义应该是:

Set<DataStream>

而不是 String。

八、你现在已经“踩过并验证”的几个关键事实

我帮你总结一下你已经实证的结论:

✅ SeaTunnel 是 DAG,不是线性 ETL ✅ 多 Source 可以汇入一个 Sink ✅ 合并是流合并,不是 SQL join ✅ Schema 对齐是前提 ✅ DSL 是数据流描述,不是 SQL

九、总结

SeaTunnel 的核心只有 3 个角色

Source     →   Transform   →   Sink
(产生流) (改流) (吃流)

数据流怎么“连”?

你只要记住这张“万能规则表”

靠的就是这两个东西:

  • plugin_output :我产生的这条数据流叫什么
  • plugin_input :我要吃哪条(或哪几条)数据流

比如,两个 source → 一个 sink

┌──────────┐
│ Source A │──┐
└──────────┘ │
├──▶ Sink
┌──────────┐ │
│ Source B │──┘
└──────────┘

一个 Source → 两个 Sink

        ┌──────▶ Sink A
Source ─┤
└──────▶ Sink B

一个 conf 里放两组“互不干扰”的流

Source A ───▶ Sink A

Source B ───▶ Sink B