- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
导读
本文是 SeaTunnel 引擎(Zeta)远程作业提交的实战手册,系统讲解如何将 SeaTunnel 作业从本地机器或 CI/CD 流水线提交到远程Zeta 集群,而非仅在本地运行。内容覆盖两种核心提交方式(--master客户端参数与 REST API)、Docker 单节点与多节点集群、Kubernetes 三种访问方式(kubectl port-forward、NodePort、LoadBalancer)、ConfigMap 作业配置管理,以及 Amazon EKS 上的 Helm 部署。读完本文,你将掌握 Zeta 集群的网络端口规划(5801 / 8080 / 8090)、REST API 提交与监控接口的完整用法,以及一套可直接复制的故障排查方法论。
1. 前置条件
向远程 Zeta 集群提交作业前,需要满足以下三项基本条件:
| 需求 | 说明 |
|---|---|
| SeaTunnel 客户端已安装 | 本地拥有 SeaTunnel 目录,可执行bin/seatunnel.sh |
| 集群可达 | 提交机器能够访问 REST API 端口(默认8080) |
| 作业配置文件就绪 | HOCON.conf、JSON 或 SQL 格式的作业配置文件 |
需要特别强调的是,这里所说的"远程集群"必须是已经启动运行的 Zeta 集群。Zeta 引擎支持三种部署形态(详见 Zeta 安装部署):本地模式(Local,仅用于测试,每个任务启动独立进程)、混合集群模式(Master 与 Worker 同进程、所有节点均可参与选举)和分离集群模式(Master 与 Worker 分离,Master 只负责作业调度、REST API 与任务提交)。其中分离集群模式是官方推荐的生产部署方式,因为 Master 不运行同步任务,负载更小、稳定性更高,即使 Worker 节点宕机也不会导致 IMap 状态数据重新分布,详见 分离集群模式部署。本文涉及的"远程提交"适用于混合集群与分离集群两种模式。
2. 从本地机器向远程集群提交
2.1--master参数(Zeta)
所有 SeaTunnel Zeta 客户端命令均支持--master参数,用于指定集群连接地址:
bin/seatunnel.sh \ --config job.conf \ --master seatunnel://192.168.1.100:5801Zeta 集群内部默认端口为5801,与 REST API 端口(8080)不同。--master参数用于通过 Hazelcast 成员协议直接连接集群。这一点可以从发行包自带的 config/hazelcast.yaml 中得到印证:hazelcast.network.port.port: 5801,且auto-increment: false(端口不会自动递增,避免多节点端口冲突)。同时该文件中join.tcp-ip.enabled: true并配置了member-list,说明 Zeta 默认使用 TCP-IP 成员发现机制;在 Kubernetes 部署时则切换到 Kubernetes 服务发现,见 deploy/kubernetes/seatunnel/conf/hazelcast-master.yaml 中的join.kubernetes.enabled: true。
Hazelcast 客户端通过 5801 端口加入集群成员组后,作业会被提交给当前 Active Master 进行调度。
2.2 使用 REST API(推荐用于自动化场景)
对于 CI/CD 流水线和脚本,推荐使用 REST API 提交,避免在构建机或调度机上维护整套 SeaTunnel 客户端:
curl -X POST http://192.168.1.100:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.jsonREST API 由 Zeta 引擎内嵌的 Jetty 服务提供。这里有两个容易混淆的"默认值"来源需要注意(详见 REST API v2 参考):
- 代码默认值:
enable-http = false、port = 8080。也就是说,如果你使用精简配置文件或删除了enable-http配置,Jetty 默认不会启动,REST API 与 Web UI 会一起不可用; - 发行包自带的
seatunnel.yaml:默认写入了enable-http: true和port: 8080。
因此,直接使用发行包自带配置启动时,REST API 通常监听http://<host>:8080/。参考 config/seatunnel.yaml 中的实际配置:
seatunnel: engine: http: enable-http: true port: 8080 enable-dynamic-port: falseREST API 仅在作业运行于 Zeta 引擎时可用;作业运行在 Flink 或 Spark 引擎上时,需要使用对应引擎自身的工具提交和监控作业。
3. Docker:单节点提交
3.1 启动已启用 REST API 的 Zeta 容器
docker run -d --name seatunnel \ -p 8080:8080 \ -e ST_DOCKER_MEMBER_COUNT=1 \ apache/seatunnel:<version>ST_DOCKER_MEMBER_COUNT是 SeaTunnel 官方 Docker 镜像约定成员数量的环境变量(单节点为 1,多节点集群需在所有节点上设置相同的期望成员数,让各容器通过 Hazelcast 发现彼此并组成集群)。
3.2 从容器外部提交作业
容器启动后,即可在宿主机上通过 REST API 提交一个FakeSource -> Console的冒烟作业:
curl -X POST http://localhost:8080/submit-job \ -H "Content-Type: application/json" \ -d '{ "env": { "job.name": "test", "job.mode": "BATCH" }, "source": [{ "plugin_name": "FakeSource", "plugin_output": "fake", "row.num": 10, "schema": { "fields": { "id": "int", "name": "string" } } }], "transform": [], "sink": [{ "plugin_name": "Console", "plugin_input": ["fake"] }] }'请求成功后会返回类似{"jobId": 733584788375666689, "jobName": "test"}的 JSON 响应,其中jobId是后续查询作业状态、停止作业的唯一标识。
3.3 在容器内执行本地 smoke 测试
docker run -d --name seatunnel \ -p 8080:8080 \ -v /path/to/your/jobs:/jobs \ apache/seatunnel:<version> # 在容器内执行 docker exec seatunnel \ /opt/seatunnel/bin/seatunnel.sh --config /jobs/my-job.conf --master local特别注意:该命令仅适用于在容器内做快速本地 smoke 测试。它不是向远程 Zeta 集群提交作业,因为--master local会在当前容器进程内本地启动作业,作业生命周期与容器进程绑定。真正测试"远程提交"应该使用第 2 节中的--master seatunnel://<host>:5801或 REST API 方式。
4. Docker:多节点集群
4.1 Docker Compose 示例
用 Docker Compose 快速搭建一个 2 节点 Zeta 集群(Master + Worker),两个容器属于同一个 bridge 网络,通过 5801 端口互相发现:
version: "3.8" services: master: image: apache/seatunnel:<version> container_name: seatunnel-master ports: - "8080:8080" - "5801:5801" environment: ST_DOCKER_MEMBER_COUNT: 2 networks: - st-net worker: image: apache/seatunnel:<version> container_name: seatunnel-worker environment: ST_DOCKER_MEMBER_COUNT: 2 networks: - st-net depends_on: - master networks: st-net: driver: bridge启动集群:
docker-compose up -d向 master 提交作业:
curl -X POST http://localhost:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.json注意只有master容器对外暴露了 5801 和 8080 两个端口,worker容器仅加入内部网络st-net,这正是"REST API 仅需在 Master 节点上对外暴露"这一网络规划原则的体现。
4.2 网络端口要求(Docker)
| 端口 | 协议 | 用途 |
|---|---|---|
| 5801 | TCP | Hazelcast 集群内部成员通信 |
| 8080 | TCP | REST API(作业提交 / 监控) |
需确保 Docker 网络中所有集群成员之间 5801 端口互通。REST API 仅需在 Master 节点上对外暴露。
5. Kubernetes:作业提交
在 Kubernetes 上部署 SeaTunnel 后(部署清单与模板见 deploy/kubernetes/seatunnel,Master 使用 StatefulSet 形态提供稳定的 Pod 名称,如seatunnel-master-0),有三种方式访问 Master 的 REST API。
5.1 使用kubectl port-forward(开发 / 临时提交)
将 Master Pod 的 REST 端口转发到本地:
# 查找 master pod kubectl get pods -n seatunnel # 转发 REST 端口 kubectl port-forward -n seatunnel \ pod/seatunnel-master-0 8080:8080在另一个终端提交作业:
curl -X POST http://localhost:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.jsonport-forward方式适合开发调试与临时验证,但存在两个天然限制:空闲超时会断开连接,Pod 重建后需重新执行转发命令。
5.2 使用 NodePort 服务(测试 / 生产)
若集群通过NodePort服务暴露 Master:
# 获取 NodePort kubectl get svc -n seatunnel seatunnel-master-rest # 使用节点 IP 和节点端口提交 curl -X POST http://<node-ip>:<node-port>/submit-job \ -H "Content-Type: application/json" \ -d @job.json5.3 使用 LoadBalancer 服务
在云厂商托管的 Kubernetes 集群(如 EKS、GKE、ACK)中,更推荐使用 LoadBalancer:
LB_IP=$(kubectl get svc -n seatunnel seatunnel-master-rest \ -o jsonpath='{.status.loadBalancer.ingress[0].ip}') curl -X POST http://${LB_IP}:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.json注意:AWS 的 LoadBalancer 返回的是
hostname而非ip,此时应改用{.status.loadBalancer.ingress[0].hostname}取值,详见第 7.2 节。
6. Kubernetes:通过 ConfigMap 管理作业配置
在集群内运行作业时,建议将作业配置文件以 ConfigMap 形式挂载,而非打包到镜像中,这样修改作业只需更新 ConfigMap 而无需重建镜像。
首先创建 ConfigMap(以 CDC 作业为例):
apiVersion: v1 kind: ConfigMap metadata: name: seatunnel-job-config namespace: seatunnel data: cdc-job.conf: | env { job.name = "cdc-prod" job.mode = STREAMING checkpoint.interval = 30000 } source { MySQL-CDC { ... } } sink { ... }在 Pod spec 中挂载:
volumeMounts: - name: job-config mountPath: /opt/seatunnel/jobs volumes: - name: job-config configMap: name: seatunnel-job-config通过kubectl exec提交:
kubectl exec -n seatunnel seatunnel-master-0 -- \ /opt/seatunnel/bin/seatunnel.sh \ --config /opt/seatunnel/jobs/cdc-job.conf这种方式无需从外部网络访问 REST API,直接在 Master Pod 内执行客户端命令即可,适用于集群网络受限(无公网/无 LoadBalancer)的场景。
7. Amazon EKS / Helm 部署
7.1 Helm 安装
SeaTunnel 官方 Helm Chart 支持一键部署 Master 与 Worker 分离架构:
helm repo add seatunnel https://apache.github.io/seatunnel-helm-charts helm repo update helm install seatunnel seatunnel/seatunnel \ --namespace seatunnel \ --create-namespace \ --set master.replicaCount=2 \ --set worker.replicaCount=4 \ --set master.service.type=LoadBalancerChart 的默认值位于 deploy/kubernetes/seatunnel/values.yaml,其中master.replicas与worker.replicas默认均为"2",且 Master 与 Worker 均配置了基于hazelcast-port(5801)的 liveness/readiness 探针;image.registry默认为apache/seatunnel,可通过image.tag指定具体版本。
7.2 EKS 获取 Load Balancer 主机名
AWS EKS 的 LoadBalancer 服务通常返回 DNS 主机名而非 IP:
kubectl get svc -n seatunnel seatunnel-master \ -o jsonpath='{.status.loadBalancer.ingress[0].hostname}'以该主机名作为 API 端点:
export ST_HOST=$(kubectl get svc -n seatunnel seatunnel-master \ -o jsonpath='{.status.loadBalancer.ingress[0].hostname}') curl -X POST http://${ST_HOST}:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.json7.3 通过 Helm values 自定义资源配置
生产环境通常需要为 Master 与 Worker 分别规划资源配额、副本数以及引擎级配置,可以编写独立的 values 文件:
# values-prod.yaml master: replicaCount: 2 resources: requests: memory: "4Gi" cpu: "2" limits: memory: "8Gi" cpu: "4" worker: replicaCount: 8 resources: requests: memory: "8Gi" cpu: "4" limits: memory: "16Gi" cpu: "8" seatunnel: config: engine: backup-count: 2 queue-type: blockingqueue print-execution-info-interval: 60 http: enable-http: true port: 8080应用配置:
helm upgrade seatunnel seatunnel/seatunnel \ --namespace seatunnel \ -f values-prod.yaml这里的seatunnel.config.engine段最终会渲染为引擎配置并注入 Pod。其中:
backup-count: 2:控制 Hazelcast IMap 状态数据的备份副本数,直接影响集群 HA 能力与 Master 节点数量规划(建议取max(1, min(5, N/2)),N 为 Master 数量);queue-type: blockingqueue:选择 Pipeline 内部队列实现;print-execution-info-interval: 60:周期性打印作业执行信息的时间间隔(秒);http.enable-http: true/http.port: 8080:启用 REST API 与 Web UI(对应 config/seatunnel.yaml 中的seatunnel.engine.http配置段)。
8. 网络与端口要求
无论采用哪种部署方式,Zeta 集群涉及三个核心端口,必须提前规划好防火墙、安全组与 Kubernetes NetworkPolicy:
| 端口 | 协议 | 使用方 | 注意事项 |
|---|---|---|---|
| 5801 | TCP | Hazelcast 集群 | 成员间通信;生产环境不要对外暴露 |
| 8080 | TCP | REST API | 生产环境建议通过认证网关暴露 |
| 8090 | TCP | Web UI | 可选,仅用于管理看板 |
在 Kubernetes 中,建议设置 NetworkPolicy 将 5801 端口限制在 SeaTunnel 命名空间内:
apiVersion: networking.k8s.io/v1 kind: NetworkPolicy metadata: name: seatunnel-internal namespace: seatunnel spec: podSelector: matchLabels: app: seatunnel ingress: - from: - namespaceSelector: matchLabels: name: seatunnel ports: - port: 5801补充说明:8080 端口上的 REST API 与 Web UI 由同一个内嵌 Jetty 服务提供(见 REST API v2 参考),只要 Jetty 未启动,两者会一起不可用。若在seatunnel.yaml中配置了enable-dynamic-port: true,实际监听端口会在port到port + port-range之间自动挑选,此时应以启动日志SeaTunnel REST service will start on port xxx为准,而不是想当然地使用 8080。
9. 提交后的监控与作业生命周期管理(REST API 进阶)
作业提交成功只是第一步。在远程集群场景下,通常需要通过 REST API 持续监控作业状态与指标。以下是实际运维中最常用的几个接口(完整参考见 REST API v2):
9.1 查询作业状态
# 运行中作业列表(支持分页) curl "http://<host>:8080/running-jobs?page=1&rows=10" # 指定作业详细信息(含 SourceReceivedCount / SinkWriteCount 等指标) curl "http://<host>:8080/job-info/<jobId>" # 已结束作业(state 取值:FINISHED / CANCELED / FAILED / SAVEPOINT_DONE / UNKNOWABLE) curl "http://<host>:8080/finished-jobs/FINISHED?page=1&rows=10"/job-info/:jobId返回的指标字段包括:SourceReceivedCount(源端接收行数)、SourceReceivedQPS、SinkWriteCount(Sink 写入尝试行数)、SinkCommittedCount(checkpoint 成功后的已提交行数)等;作业运行中还可返回diagnostics字段,其中pipelines[].restoreCount如果持续增长而jobStatus一直是RUNNING,说明作业正处于崩溃重启循环,需要结合日志排查。
9.2 集群与资源概览
# 集群概览(totalSlot / runningJobs / pendingJobs 等) curl "http://<host>:8080/overview" # Worker 资源快照(totalSlots / freeSlots / cpuUsage / memUsage) curl "http://<host>:8080/resource/workers" # Pending 队列诊断(排查作业长时间 WAITING / PENDING 的原因) curl "http://<host>:8080/pending-jobs?limit=10&pretty=true"/pending-jobs是排查资源不足的利器:当作业提交后长时间处于PENDING,响应中的lackingTaskGroups、failureMessage(如NoEnoughResourceException: slot not enough)和blockingJobIds可以直接指明缺多少 Slot、被哪些作业占用。
9.3 停止作业
curl -X POST http://<host>:8080/stop-job \ -H "Content-Type: application/json" \ -d '{"jobId": 733584788375666689, "isStopWithSavePoint": false, "force": false}'参数说明:
| 参数 | 是否必传 | 说明 |
|---|---|---|
jobId | 是 | 作业 ID |
isStopWithSavePoint | 否 | 是否通过 savepoint 方式停止(保存当前状态,便于后续恢复) |
force | 否 | 是否强制停止(忽略isStopWithSavePoint),仅应在异常场景使用,因为可能导致检查点数据不完整 |
9.4 上传配置文件提交
除 JSON 请求体外,REST API 还支持直接上传.conf(HOCON)、.sql、.json文件:
curl --location 'http://127.0.0.1:8080/submit-job/upload' \ --form 'config_file=@"/temp/fake_to_console.conf"'上传大小受seatunnel.engine.http.upload-max-file-size-mb(默认 10 MB)与upload-max-request-size-mb(默认 10 MB)限制,超出会在解析配置之前被拒绝。
注意:REST API 不支持 dry-run(仅 CLI 提供);
/submit-job还支持format参数(json/hocon/sql,默认json),以及基于restoreMode/restoreSourceJobId/isStartWithSavePoint的作业恢复能力。
10. 故障排查
远程提交场景中,问题往往出在"网络不通"、"端口未开"或"资源不足"三类。下表总结了常见现象、可能原因与修复方法:
| 现象 | 可能原因 | 修复方法 |
|---|---|---|
8080 端口Connection refused | REST API 未启用或端口错误 | 设置enable-http: true;检查port配置 |
| Worker 无法加入集群 | 防火墙阻断 5801 端口 | 开放所有集群节点间 TCP 5801 |
kubectl port-forward断开 | 空闲超时或 Pod 重启 | 重新执行 port-forward;考虑改用 NodePort |
作业已提交但状态始终为WAITING | 无可用 Worker 槽位 | 扩容 Worker 副本数或检查资源配额 |
| EKS LoadBalancer 主机名无法解析 | DNS 传播延迟 | 等待 1–2 分钟;用nslookup验证 |
Helm 安装卡在pending-install | 上次安装失败残留 | 执行helm rollback或helm uninstall后重试 |
补充两个容易被忽略的排查点:
- 8080 打不开时先确认 Jetty 是否真的启动:
seatunnel.engine.http.enable-http(或enable-https)才是 REST API / Web UI 的开关,仅配置hazelcast.yaml中的network.rest-api.enabled不能替代 Jetty 开关(详见 REST API v2); - 作业状态
WAITING/PENDING时用/pending-jobs拿诊断信息:响应中的failureMessage会直接告诉你是否NoEnoughResourceException,以及具体缺哪些 TaskGroup 的 Slot。
参考
- REST API v2 完整参考
- Zeta 引擎安装部署
- 分离集群模式部署
- SeaTunnel 引擎配置文件(示例)
- Hazelcast 集群配置文件(示例)
- Kubernetes Helm Chart 默认值
- Kubernetes 部署 Master Hazelcast 配置
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
向远程 Zeta 集群提交 SeaTunnel 作业:从 Docker 单机到 Kubernetes/EKS 的完整实战指南
向远程 Zeta 集群提交 SeaTunnel 作业:从 Docker 单机到 Kubernetes/EKS 的完整实战指南 本文面向需要将 SeaTunnel
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Zeta 引擎 RESTful API V2 完全指南:监控、作业提交与集群运维
SeaTunnel Zeta 引擎 RESTful API V2 完全指南:监控、作业提交与集群运维 SeaTunnel(Zeta 引擎)内置了一套基于 HTT
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel 集群 Helm 部署实战:从 Chart 安装到任务提交完整指南
SeaTunnel 集群 Helm 部署实战:从 Chart 安装到任务提交完整指南 SeaTunnel(Apache SeaTunnel)是一个多模态、高性能
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考