ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

SeaTunnel 入门指南:引擎选型、首次作业运行与学习路径

SeaTunnel 入门指南:引擎选型、首次作业运行与学习路径 SeaTunnel 入门指南引擎选型、首次作业运行与学习路径【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇指南是 SeaTunnel 开源数据集成平台的官方入门总览对应仓库文档 docs/en/getting-started/overview.md面向首次接触 SeaTunnel 的开发者。文章将带你完成三件事理解 SeaTunnel 的核心应用场景与统一连接器模型、在 Zeta / Flink / Spark 三种执行引擎之间做出正确选型、以及从部署安装到跑通第一个FakeSource - FieldMapper - Console示例作业的完整实操路径并在此基础上规划后续的进阶学习路线。SeaTunnel 能帮你做什么SeaTunnel 是一个面向异构系统间数据搬运而设计的分布式数据集成平台其核心在于统一的连接器模型无论数据源是数据库、文件、数据仓库、消息队列还是对象存储都可以用同一套source / transform / sink的作业描述方式完成读写。在实际业务中团队通常用它解决以下四类问题数据库、文件与数据仓库之间的批量同步离线全量/增量 ETL例如把 MySQL 中的数据定时搬运到 ClickHouse、Doris 或数据湖CDC 与实时同步通过 CDC 连接器捕获数据库变更事件实时同步到下游系统多表或全库迁移一条作业描述多个表甚至整个数据库的迁移任务配合引擎级的多表同步能力多模态数据摄入覆盖结构化、非结构化乃至二进制数据的读取与写入。如果你正在评估 SeaTunnel 是否适合你的项目建议直接以内置的SeaTunnel EngineZeta作为起点它是新部署的默认引擎配置路径最短适合快速在本地验证方案。从 docs/en/introduction/about.md 的定位描述可以看到SeaTunnel 的设计哲学是一个作业本质上就是一条管道你在配置文件中描述作业SeaTunnel 将其运行成一条从Source到Transform再到Sink的数据管道连接器定义你读写什么引擎则决定作业在哪里运行。这一连接器优先 引擎可插拔的设计使 Source、Transform、Sink 插件可以在不同引擎之间复用。选择执行引擎Zeta、Flink 还是 SparkSeaTunnel 支持多种执行引擎选择哪一个直接决定了你的部署形态和运维方式。总体选型规则如下SeaTunnel EngineZeta如果你没有现成的 Flink / Spark 基础设施从它开始。它是为数据集成场景原生构建的引擎也是大多数新部署的默认推荐Apache Flink如果你的团队已经运维 Flink 集群希望复用现有运行栈Apache Spark如果你已经运行 Spark且工作负载以批量为主。三种引擎快速对比引擎描述推荐场景SeaTunnel Engine (Zeta)专为数据集成而生的原生引擎新项目、数据同步Apache Flink分布式流处理引擎已有 Flink 基础设施Apache Spark分布式批/流处理引擎已有 Spark 基础设施特性对比特性SeaTunnel EngineFlinkSpark批量处理Batch✅✅✅流式处理Streaming✅✅✅CDC 支持✅✅❌精确一次Exactly-Once✅✅✅多表同步✅✅✅Schema 演进✅✅❌REST API✅❌❌Web UI✅✅✅单机模式✅✅✅集群模式✅✅✅说明上表中的REST API特指 SeaTunnel 自身的作业提交/监控 API详见 docs/en/engines/zeta/rest-api-v2.md仅由 SeaTunnel EngineZeta服务端实现因此当作业运行在 Flink 或 Spark 引擎上时不可用。使用 Flink / Spark 引擎时需要通过对应引擎自身的工具链提交与监控作业如 Flink CLI/REST API、Spark 的spark-submit/History Server。使用体验对比维度SeaTunnel EngineFlinkSpark吞吐量⭐⭐⭐ 高⭐⭐ 中⭐⭐ 中延迟⭐⭐⭐ 低⭐⭐⭐ 低⭐⭐ 中资源占用⭐⭐⭐ 低⭐⭐ 中⭐ 高启动时间⭐⭐⭐ 快⭐⭐ 中⭐ 慢安装难度⭐⭐⭐ 简单⭐⭐ 中⭐⭐ 中配置复杂度⭐⭐⭐ 简单⭐⭐ 中⭐⭐ 中外部依赖⭐⭐⭐ 无⭐⭐ Zookeeper可选⭐ YARN/Mesos学习曲线⭐⭐⭐ 平缓⭐⭐ 中等⭐⭐ 中等完整的特性与性能对比可参考 docs/en/engines/overview.md。各引擎的适用场景SeaTunnel EngineZeta—— 默认推荐适合全新的数据集成项目数据同步与 CDC 场景没有现成大数据基础设施的团队需要低资源消耗的场景例如大量小表的实时同步数据库迁移类项目。其核心优势包括无外部依赖不需要 Zookeeper、HDFS 即可完成集群管理与 HA、针对数据同步场景深度优化动态线程共享、JDBC 连接复用、Pipeline 级容错、内置集群管理与高可用。典型用例是MySQL 到 ClickHouse 实时同步、多表 CDC 同步和数据库迁移。Apache Flink适合已有 Flink 基础设施、需要复杂流处理能力或希望深度整合 Flink 生态如 Flink SQL、高级状态管理的团队Apache Spark适合已有 Spark 基础设施、以大规模批量处理为主的团队如与 Hive/HDFS 集成、YARN/Kubernetes 部署、MLlib/GraphX 生态。关于引擎层面的实现细节可进一步阅读 docs/en/engines/zeta/about.mdZeta 引擎以Pipeline 作为 checkpoint 与容错的最小粒度单个任务失败只会影响其上下游任务避免整个作业失败或回滚通过动态线程共享技术在实时同步大量小数据量表时共享线程、减少不必要的线程创建在 CDC 场景复用日志读取与解析资源并在读写两侧尽量减少 JDBC 连接数量。集群管理上支持单机、集群以及**自主集群去中心化**模式——相同cluster_name的节点会自动组网并在主节点故障时自动选出新主节点。使用 SeaTunnel 之前需要准备什么在开始部署之前请确认以下前置条件Java 8 或 Java 11并正确配置JAVA_HOME更高版本理论上也可工作见 docs/en/getting-started/locally/deployment.md一份 SeaTunnel 二进制发布包从官方下载页获取seatunnel-version-bin.tar.gz所需的连接器插件安装在${SEATUNNEL_HOME}/connectors/目录下所选连接器要求的第三方驱动 JAR例如 MySQL 的 JDBC 驱动。如果只想跑通示例作业通常只需要connector-fake和connector-console两个插件即可。这一点与 docs/en/getting-started/job-configuration-guide.md 中Validation Checklist给出的运行前检查项一致Java 与JAVA_HOME正确、必备连接器已安装、第三方驱动就位、源端凭据与网络可达、目标表/主题/路径已存在如要求、job.mode与所用连接器能力匹配。推荐首次运行最短路径跑通第一个作业如果你想用最短路径验证安装是否成功按以下顺序操作阅读 部署文档 并安装二进制包安装首个作业所需的示例插件以FakeSource - FieldMapper - Console跑通本地 SeaTunnel Engine 快速开始示例成功后将演示用 source 和 sink 替换为真实连接器。第一步部署 SeaTunnel 与连接器下载二进制包Linux/macOSexport version3.0.0 wget https://archive.apache.org/dist/seatunnel/${version}/apache-seatunnel-${version}-bin.tar.gz tar -xzvf apache-seatunnel-${version}-bin.tar.gz自 2.2.0-beta 起二进制包默认不再附带连接器依赖首次使用需要运行以下命令安装连接器也可以从 Apache Maven 仓库手动下载连接器放入connectors/目录2.3.5 之前的版本需放入connectors/seatunnel目录sh bin/install-plugin.shWindows 下使用打包内批处理脚本基于 Maven Wrapper无需单独安装 Mavencd apache-seatunnel-3.0.0 bin\install-plugin.cmd如需安装特定发布版本的连接器可在命令后附带版本号例如sh bin/install-plugin.sh 3.0.0。对于已发布的连接器版本Linux/macOS 上的install-plugin.sh会通过 HTTPS 直接下载 JAR 与校验和因此无需 Maven该路径要求本机具备curl、mktemp以及sha512sum/sha1sum/shasum/openssl中的一种。Windows 的install-plugin.cmd仍使用内置 Maven Wrapper。若需走 HTTPS Maven 镜像可通过环境变量指定SEATUNNEL_MAVEN_REPOSITORYhttps://repo.example.com/maven2 \ sh bin/install-plugin.sh 3.0.0通常你不需要安装全部连接器。通过修改config/plugin_config仓库中的实际示例见 config/plugin_config只启用所需插件。例如让示例作业可用只需保留--seatunnel-connectors-- connector-fake connector-console --end--所有支持的连接器及其在plugin_config中的配置名称可以在${SEATUNNEL_HOME}/connectors/plugins-mapping.properties中查到。在仓库源码侧connector-fake的实现在 seatunnel-connectors-v2/connector-fake/src/main/java/org/apache/seatunnel/connectors/seatunnel/fake/source/FakeSource.javaconnector-console的实现在 seatunnel-connectors-v2/connector-console/src/main/java/org/apache/seatunnel/connectors/seatunnel/console/sink/ConsoleSink.java读者可以按需深入阅读。第二步编写作业配置文件编辑config/v2.batch.config.template仓库中的真实模板见 config/v2.batch.config.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 } }关于配置概念更完整的说明见 配置基础概念。第三步运行 SeaTunnel 应用cd apache-seatunnel-${version} ./bin/seatunnel.sh --config ./config/v2.batch.config.template -m localWindows 下运行等价批处理入口cd apache-seatunnel-3.0.0 bin\seatunnel.cmd --config config\v2.batch.config.template -m local提示自 2.3.1 起seatunnel.sh中的参数-e已废弃请改用-m。查看输出命令运行后SeaTunnel 控制台会打印类似如下的日志这是判断命令是否成功的标志。ConsoleSinkWriter会输出转换后的行类型new_nameSTRING, ageINT以及逐行写入的数据rowIndex1到16与row.num 16对应2022-12-19 11:01:45,417 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - output rowType: new_nameSTRING, ageINT 2022-12-19 11:01:46,489 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex1: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: CpiOd, 8520946 2022-12-19 11:01:46,490 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex2: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: eQqTs, 1256802974 ... 2022-12-19 11:01:46,491 INFO org.apache.seatunnel.connectors.seatunnel.console.sink.ConsoleSinkWriter - subtaskIndex0 rowIndex16: SeaTunnelRow#tableId-1 SeaTunnelRow#kindINSERT: mIJDt, 995616438从实现看该输出由 ConsoleSinkWriter.java 负责打印其单元测试 ConsoleSinkWriterTest.java 验证了 writer 的行为。理解作业配置env / source / transform / sinkSeaTunnel 作业是声明式定义的绝大多数集成无需编写代码只需在配置文件中描述执行环境、源端、可选转换和汇端。绝大多数作业遵循相同的顶层结构核心是四个块env控制作业如何执行source定义数据从哪来transform在途修改数据可选sink定义数据去哪。env块中的常用公共设置如下配置键含义job.modeBATCH或STREAMINGparallelism作业默认并行度job.name可选的作业显示名称checkpoint.interval流式作业与精确一次工作流的 checkpoint 间隔若使用 Flink 或 Spark 引擎引擎专属参数也在env中配置见 JobEnvConfig。仓库模板 config/v2.batch.config.template 给出的实际示例为parallelism 2、job.mode BATCH、checkpoint.interval 10000。source块通常包含连接器名称、连接参数、读取范围表、topic、路径或查询以及 schema/格式相关参数transform块用于字段重命名/映射、行过滤、行类型转换、SQL 转换或写入前校验——如果源与目标 schema 直接对齐完全可以省略 transform 块直接从 source 到 sinksink块则包含连接器名称、连接参数、目标表/topic/路径、写入语义或批处理参数。理解 plugin_input 与 plugin_output这两个键是理解数据如何在作业中流动的最重要约定plugin_output为 source 或 transform 产出的数据流命名plugin_input告诉 transform 或 sink 消费哪条上游数据流。当一条作业读取多个 source、一个 transform 扇出到多个 sink或作业存在多个分支时显式命名能让管道结构清晰可维护。如果作业只有一条上游路径SeaTunnel 通常可以按默认约定工作无需两个字段都写但出于可读性仍建议显式命名。支持的配置格式SeaTunnel 支持多种配置风格HOCON默认且最常用的格式JSON适合由其他系统生成配置的场景SQL适合以 SQL 为中心的工作流。格式细节可参考 配置概念 与 SQL 配置。进阶实战MySQL 到 Doris 的批量同步当示例作业跑通后最自然的下一步是把演示插件替换为真实连接器。以 MySQL 到 Doris 的批量同步为例Step 1下载连接器。在${SEATUNNEL_HOME}/config/plugin_config中加入连接器名称然后执行安装命令确保connector-jdbc与connector-doris出现在${SEATUNNEL_HOME}/connectors/目录也可以从 Apache Maven 仓库手动下载放入该目录--seatunnel-connectors-- connector-jdbc connector-doris --end--sh bin/install-plugin.shStep 2放置 MySQL 驱动。下载 MySQL JDBC 驱动 JAR放入${SEATUNNEL_HOME}/lib/目录。Step 3编写作业配置。例如保存为seatunnel/job/st.confenv { parallelism 2 job.mode BATCH } source { Jdbc { url jdbc:mysql://localhost:3306/test driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 user user password pwd table_path test.table_name query select * from test.table_name } } sink { Doris { fenodes doris_ip:8030 username user password pwd database test_db table table_name sink.enable-2pc true sink.label-prefix test-cdc doris.config { format json read_json_by_linetrue } } }Step 4运行作业cd seatunnel/ ./bin/seatunnel.sh --config ./job/st.conf -m local运行结束后控制台会打印类似如下的作业统计信息用于核对读写数量与失败数量*********************************************** Job Statistic Information *********************************************** Start Time : 2024-08-13 10:21:49 End Time : 2024-08-13 10:21:53 Total Time(s) : 4 Total Read Count : 1000 Total Write Count : 1000 Total Failed Count : 0 ***********************************************如需进一步调优作业可参考 Source-MySQL 与 Sink-Doris 连接器文档。其余端到端配方source 到 sink 的完整走通示例还包括 MySQL CDC to Kafka、MySQL CDC to Doris、JDBC to S3、Kafka to Iceberg、Http to JDBC、File to StarRocks 与 Multi-table CDC。从示例到真实作业替换插件的通用方法把示例作业改造成真实作业的最快路径是渐进式替换保留示例中的env块将FakeSource替换为真实 source 连接器将Console替换为目标 sink 连接器仅当源 schema 与目标 schema 无法直接对齐时添加 transform按需补充连接器专属 JAR 或驱动。常见的真实作业形态包括MySQL 到 Doris、Kafka 到 Iceberg、S3File 到 StarRocks、PostgreSQL CDC 到 Kafka。推荐阅读路径路径 A我只想跑通第一个作业部署文档SeaTunnel Engine 快速开始作业配置指南路径 B我已经明确要构建的作业作业配置指南Source 连接器Sink 连接器Transforms场景配方路径 C我需要先理解架构关于 SeaTunnel它是如何工作的架构总览首次运行成功之后一旦示例作业跑通下一步通常从以下三者中选择其一将FakeSource与Console替换为真实连接器从本地验证切换到集群部署暴露 REST API 与 Web UI 以获得运维可见性。对应可继续阅读作业配置指南场景配方SeaTunnel EngineZeta部署REST API 与 Web UI向远程 Zeta 集群提交作业小结从总览到第一条生产管道SeaTunnel 的价值在于用一套统一连接器模型覆盖批同步、实时同步、CDC 与多表迁移等绝大部分数据集成场景并通过可插拔的引擎设计让你从零基础设施的 Zeta 单机验证平滑演进到 Flink / Spark 集群或 Zeta 多节点集群。本文给出的最短路径——部署二进制包、安装connector-fake与connector-console、以FakeSource - FieldMapper - Console跑通本地作业、再逐步替换为真实连接器——是进入 SeaTunnel 生态最直接的一扇门后续的配置结构env / source / transform / sink与plugin_input / plugin_output、连接器文档与场景配方则构成了从入门到生产落地完整的知识闭环。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表