SeaTunnel 基于 Flink 引擎的本地快速上手:部署、配置与运行实战
2026/9/18 9:22:59 网站建设 项目流程

SeaTunnel 基于 Flink 引擎的本地快速上手:部署、配置与运行实战

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

导读

本指南面向已经拥有或计划使用 Apache Flink 运行环境的团队,讲解如何在 SeaTunnel 中把 Flink 作为执行引擎,从零跑通一个完整的「Flink 本地快速开始」示例。你将掌握 Flink 与 SeaTunnel 的部署衔接、FLINK_HOME环境配置、HOCON 格式作业配置文件的编写(Source / Transform / Sink 三段式结构)、按 Flink 版本选择合适的启动脚本,并理解 SeaTunnel 是如何通过翻译层把自身连接器 API 适配到 Flink 运行时之上的。若你是首次评估 SeaTunnel 且没有既有的 Flink 集群,建议先走内置 Zeta 引擎的 Quick Start With SeaTunnel Engine,只有当你确实需要 Flink 时再回到本指南。

阅读前置:理解 SeaTunnel 与 Flink 的关系

SeaTunnel 本身是一套独立的数据集成工具,拥有内置引擎 Zeta;但它同样支持把 Flink 作为外部执行引擎来运行同步作业。选择 Flink 通常出于以下场景:

  • 团队已经在生产环境运维 Flink 集群,希望复用现有的 Flink 部署、监控与运维体系;
  • 同步作业需要与更大的 Flink 流处理环境对齐;
  • 需要借助 Flink 成熟的 checkpoint 语义与状态管理能力。

在深入本指南之前,建议按顺序阅读以下文档建立整体认知:

  • Engine Overview:了解 SeaTunnel 支持的多引擎架构;
  • SeaTunnel With Flink:了解 Flink 引擎的适用场景与专属配置;
  • Job Configuration Guide:了解作业配置文件的通用结构。

第 1 步:部署 SeaTunnel 与连接器插件

在配置 Flink 之前,需要先完成 SeaTunnel 本体及其连接器插件的部署,完整步骤见 Deployment。核心要点如下:

  1. 安装 Java 并设置JAVA_HOME:SeaTunnel 需要 Java 8 或 Java 11(理论上高于 Java 8 的版本均可运行)。
  2. 下载二进制发行包:从 Apache 官方发行渠道获取apache-seatunnel-<version>-bin.tar.gz(Windows 对应.zip包),解压后得到${SEATUNNEL_HOME}目录。
  3. 安装连接器插件:从 2.2.0-beta 起,二进制包默认不再附带连接器依赖,首次使用需执行安装脚本:
sh bin/install-plugin.sh

install-plugin.sh会直接通过 HTTPS 下载连接器 JAR 及其校验和(需要curlmktemp以及sha512sum/sha1sum/shasum/openssl之一)。如果想为本例作业最小化安装,只需保证connector-fakeconnector-console两个插件可用,可通过编辑 config/plugin_config 指定要安装的插件:

--seatunnel-connectors-- connector-fake connector-console --end--

所有受支持连接器与plugin_config中的配置名对应关系,可在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties中查到。

第 2 步:部署 Flink 并配置 SeaTunnel

下载并部署 Flink

本指南要求 Flink 版本不低于 1.12.0。请前往 Apache Flink 官方下载页获取对应版本并解压到本地目录(例如/opt/flink)。若需要了解 Flink 自身的 Standalone 集群部署方式,可参考 Flink 官方文档中「Getting Started: Standalone」部分(以 release-1.14 文档为例)。

配置 SeaTunnel 环境变量

编辑${SEATUNNEL_HOME}/config/seatunnel-env.sh,将FLINK_HOME指向 Flink 部署目录。仓库中的 config/seatunnel-env.sh 展示了该文件的结构:

# Home directory of spark distribution. SPARK_HOME=${SPARK_HOME:-/opt/spark} # Home directory of flink distribution. FLINK_HOME=${FLINK_HOME:-/opt/flink} # Whether to enable metalake (true/false). METALAKE_ENABLED=${METALAKE_ENABLED:-false} # Type of metalake implementation. METALAKE_TYPE=${METALAKE_TYPE:-gravitino} # Metalake service URL, format: http://host:port/api/metalakes/{metalake_name}/catalogs/. METALAKE_URL=${METALAKE_URL:-http://localhost:8090/api/metalakes/default_metalake_name/catalogs/}

可以看到,FLINK_HOME默认值为/opt/flink,同时支持通过同名环境变量覆盖,即export FLINK_HOME=/your/flink/path后再启动作业同样生效。该文件会被各启动脚本(例如start-seatunnel-flink-13-connector-v2.sh)在运行前 source 加载,因此修改后无需重新编译。

需要说明的是:使用 Flink 运行 SeaTunnel 同步任务时,无需部署 SeaTunnel Engine(Zeta)服务集群,SeaTunnel 只会作为作业提交方把任务交给 Flink 运行时执行。

第 3 步:编写作业配置文件

SeaTunnel 采用声明式作业定义:无需为大多数集成编写代码,只需在配置文件中描述执行环境(env)、数据源(source)、可选的转换(transform)与数据目标(sink)。

编辑config/v2.streaming.conf.template,它决定了 SeaTunnel 启动后的数据输入、处理与输出方式。仓库中的示例配置内容如下(与下方示例应用一一对应):

env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { plugin_output = "fake" row.num = 16 schema = { fields { name = "string" age = "int" } } } } transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" field_mapper = { age = age name = new_name } } } sink { Console { plugin_input = "fake1" } }

对这四个配置块的说明:

  • env:控制作业的执行方式。parallelism = 1指定作业默认并行度;job.mode = "BATCH"声明批处理模式(可选值为BATCHSTREAMING)。此外还常配置job.name(作业显示名)与checkpoint.interval(checkpoint 间隔)。
  • source:定义数据来源。示例使用FakeSource模拟数据源,row.num = 16表示生成 16 行数据,schema.fields声明了两列:name(string)与age(int)。plugin_output = "fake"为输出流命名,供下游插件引用。
  • transform:可选的中间处理环节。示例中的FieldMapper通过field_mapper完成字段映射/重命名:把name改名为new_nameage保持原名。plugin_input = "fake"消费上游输出,plugin_output = "fake1"定义新的输出流名。
  • sink:定义数据去向。示例使用Console把数据打印到控制台,plugin_input = "fake1"指向 transform 的输出。

plugin_inputplugin_output的数据流约定

这两个键是理解 SeaTunnel 数据如何在作业内部流动的最重要约定:

  • plugin_output为 source 或 transform 产生的数据流命名;
  • plugin_input告诉 transform 或 sink 消费哪条上游数据流。

当作业存在多个 source、单个 transform 扇出到多个 sink、或作业的不同分支需要保持清晰时,显式命名尤其有价值。如果作业只有单条上游链路,SeaTunnel 通常可以按默认约定省略这两个字段,但为了可读性仍推荐显式声明。

从示例走向真实作业

最快的改造路径是渐进式替换示例插件:保留env块 → 用真实 source 连接器替换FakeSource→ 用目标 sink 连接器替换Console→ 仅在源 schema 与目标 schema 不对齐时添加 transform → 按需补充连接器专属 JAR 或驱动。

更详细的配置概念可参考 Config Concept 与 Job Configuration Guide;transform 参数细节见 Transform Common Options。

第 4 步:运行 SeaTunnel 作业

根据 Flink 版本选择启动脚本

SeaTunnel 针对不同 Flink 大版本提供了独立的 starter 模块(源码位于 seatunnel-flink-starter 目录),因此启动命令需要按 Flink 版本区分。

Flink 版本在 1.12.x 与 1.14.x 之间:

cd "apache-seatunnel-${version}" ./bin/start-seatunnel-flink-13-connector-v2.sh --config ./config/v2.streaming.conf.template

Flink 版本在 1.15.x 与 1.18.x 之间:

cd "apache-seatunnel-${version}" ./bin/start-seatunnel-flink-15-connector-v2.sh --config ./config/v2.streaming.conf.template

仓库中还提供了start-seatunnel-flink-20-connector-v2.sh(对应 Flink 2.0)以及各脚本的.cmdWindows 版本,例如 start-seatunnel-flink-13-connector-v2.sh。

启动脚本的底层逻辑

以 1.13 的启动脚本为例(start-seatunnel-flink-13-connector-v2.sh),它实际执行的是:

  1. 加载config/seatunnel-env.sh,使FLINK_HOME等环境变量生效;
  2. 首次运行且存在connectorslibplugins目录时,将它们打包为runtime.tar.gz
  3. 通过java -cp ... org.apache.seatunnel.core.starter.flink.FlinkStarter解析参数并生成最终的 Flink 作业执行命令;
  4. 由脚本eval执行该命令,把作业提交到 Flink 运行时。

其中入口类FlinkStarter(见 FlinkStarter.java)的作用是「生成最终的 Flink job 执行命令」——它在main方法中调用buildCommands()并打印出完整命令串,这正是脚本能够先解析再eval提交的原因。

查看运行输出

命令运行成功后,SeaTunnel 控制台会打印类似如下的日志,可用于判断作业是否成功执行:

fields : name, age types : STRING, INT row=1 : elWaB, 1984352560 row=2 : uAtnp, 762961563 row=3 : TQEIB, 2042675010 row=4 : DcFjo, 593971283 row=5 : SenEb, 2099913608 row=6 : DHjkg, 1928005856 row=7 : eScCM, 526029657 row=8 : sgOeE, 600878991 row=9 : gwdvw, 1951126920 row=10 : nSiKE, 488708928 row=11 : xubpl, 1420202810 row=12 : rHZqb, 331185742 row=13 : rciGD, 1112878259 row=14 : qLhdI, 1457046294 row=15 : ZTkRx, 1240668386 row=16 : SGZCr, 94186144

日志先是fieldstypes两行声明了输出 schema(name, age/STRING, INT),随后逐行打印 16 条随机生成的数据记录——每条记录对应FakeSource的一行数据,并已通过FieldMapper完成字段映射。

Flink 专属配置:flink.前缀与 env 内嵌配置

env块内,SeaTunnel 作业级 Flink 配置统一使用flink.前缀。以 SeaTunnel With Flink 中的示例为例:

env { parallelism = 1 flink.execution.checkpointing.unaligned.enabled = true }

需要注意的限制:SeaTunnel 作业配置对内联枚举类型的支持并不完整,需要枚举类取值的设置(超出受支持内联类型范围的)应在 Flink 自身中配置。常见的受支持内联值类型为IntegerBooleanStringDuration

一个更完整的 Flink 最小示例作业

同样来自 SeaTunnel With Flink,下面的示例在 Flink 上运行并支持 checkpoint 配置,schema 覆盖了更丰富的字段类型:

env { parallelism = 1 checkpoint.interval = 5000 flink.execution.checkpointing.mode = "EXACTLY_ONCE" flink.execution.checkpointing.timeout = 600000 } source { FakeSource { row.num = 16 plugin_output = "fake_table" schema = { fields { c_map = "map<string, string>" c_array = "array<int>" c_string = string c_boolean = boolean c_int = int c_bigint = bigint c_double = double c_bytes = bytes c_date = date c_decimal = "decimal(33, 18)" c_timestamp = timestamp } } } } transform { FieldMapper { plugin_input = "fake_table" plugin_output = "fake_output" field_mapper = { c_string = c_string c_int = c_int } } } sink { Console { plugin_input = "fake_output" } }

env 配置在源码中的落地方式

env块中的 checkpoint 相关配置最终由 Flink 运行时环境类解析并应用到 Flink 的StreamExecutionEnvironment。以 AbstractFlinkRuntimeEnvironment.java 为例,其setCheckpoint()方法展示了关键逻辑:

  • checkpoint.interval未配置或小于等于 0 时,批处理作业直接禁用 checkpoint;
  • 流式作业在 interval 非正数时会回退到默认值 10 秒(DEFAULT_CHECKPOINT_INTERVAL_MS = 10000L);
  • checkpoint.mode支持exactly-onceat-least-once两种取值,映射到 Flink 的CheckpointingMode
  • 通过checkpoint.data.uri可设置FsStateBackend,配合state.backend = "rocksdb"可切换为RocksDBStateBackend

而 FlinkRuntimeEnvironment.java 则负责在prepare()阶段创建流执行环境,并在配置了job.name时设置作业名。

深入原理:SeaTunnel API 如何适配到 Flink

SeaTunnel 连接器开发者只需要实现引擎无关的SeaTunnelSourceSeaTunnelSinkSeaTunnelTransform接口,而 Flink 通过自己的一套运行时契约(checkpoint 生命周期、source reader 模型、sink 接口)执行作业。两者之间的桥接由 Flink 翻译层完成,详见 Flink Translation Layer。

高层映射关系

SeaTunnelSource -> FlinkSource adapter -> Flink Source runtime SeaTunnelSink -> FlinkSink adapter -> Flink Sink runtime SeaTunnel types -> serializer and type adapters -> Flink state and records

翻译层主要适配四个方面:生命周期(lifecycle)、上下文(context)、序列化(serialization)与 checkpoint 语义。

源码侧与 Sink 侧映射

  • Source 侧:把 SeaTunnel 的 reader/enumerator 模型桥接到 Flink 的 source 运行时,包括映射 boundedness、从 SeaTunnel reader 创建SourceReader适配器、从 split enumerator 创建 enumerator 适配器、包装 split 与 enumerator 的状态序列化器用于 Flink checkpoint。
  • Sink 侧:把 SeaTunnel sink 契约适配为 Flink 的 writer/committer 模型,暴露 committer 与 aggregated committer 行为,映射 writer 状态与 commit info 序列化器——对于依赖 checkpoint 驱动提交语义的 sink 尤为重要。

checkpoint 与状态对齐

Flink 是 SeaTunnel source/sink API 如此设计的主要原因之一。翻译层必须保证状态快照时机、checkpoint 完成回调、split 与 writer 状态序列化、提交协调语义的一致性。如果对齐出错,通常表现为数据重复、恢复后丢数据、checkpoint 失败或 sink 提交不一致。

相关实现位于seatunnel-translation/seatunnel-translation-flink/目录,常用类包括FlinkSourceFlinkSourceReaderFlinkSourceEnumeratorFlinkSourceReaderContext等。

常见问题定位与运行前检查清单

翻译层出现问题时,往往集中在 checkpoint 回调、序列化器兼容性、watermark/事件时间预期、以及引擎专属配置泄漏进连接器代码等环节。排查时需要区分三类问题来源:连接器 bug、SeaTunnel API 契约问题、Flink 翻译层问题。

运行作业前的验证清单:

  • Java 与JAVA_HOME已正确设置;
  • 所需连接器插件已安装(本例为connector-fakeconnector-console);
  • 需要的第三方驱动已就位;
  • source 的凭据与网络访问有效;
  • 目标表、topic 或路径(如需)已提前创建;
  • job.mode与你计划使用的连接器能力匹配。

继续深入的方向

  • 现在可以编写自己的配置文件了:根据 Source Connectors 选择你需要的连接器,并按对应文档配置参数;
  • 想全面了解 SeaTunnel 与 Flink 的结合方式,见 SeaTunnel With Flink;
  • 想理解 SeaTunnel API 如何适配 Flink 运行时,见 Flink Translation Layer;
  • SeaTunnel 内置引擎 Zeta 是默认引擎,若想要最短的本地验证路径,可参考 Quick Start With SeaTunnel Engine;从源码树运行示例时,示例模块为seatunnel-examples/seatunnel-flink-connector-v2-example,入口类为org.apache.seatunnel.example.flink.v2.SeaTunnelApiExample

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询