SeaTunnel Zeta 引擎 REST API 任务全生命周期管理实战指南
2026/9/17 21:30:56 网站建设 项目流程

SeaTunnel Zeta 引擎 REST API 任务全生命周期管理实战指南

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

本指南是 REST API v2 参考文档 的实战配套教程,聚焦于 Apache SeaTunnel(Zeta 引擎)下作业的提交、状态查询、日志获取、停止/取消/保存点、恢复重启、认证与性能调优等完整生命周期操作。读者学完后,可熟练使用 curl 与 Zeta 引擎内置 REST 服务完成一套可复制、可上线的作业管理流程,并理解底层JobInfoService、Jetty 内嵌服务等实现原理。

适用前提:本 API 由 SeaTunnel Engine(Zeta)内嵌的 Jetty 服务提供,仅对运行在Zeta 引擎上的作业生效;若作业运行在 Flink 或 Spark 引擎上,请使用对应引擎自身的提交与监控工具。


1. 前置条件:开启 REST 服务

REST API 与 Web UI 共用 Zeta 内嵌的 Jetty 服务,Jetty 仅在seatunnel.engine.http.enable-http = true(或enable-https = true)时启动。在config/seatunnel.yaml中开启:

seatunnel: engine: http: enable-http: true port: 8080 enable-dynamic-port: true port-range: 100
  • enable-http:是否启动 HTTP 服务(代码默认false,打包的seatunnel.yaml示例已默认开启);
  • port:固定监听端口(默认8080);
  • enable-dynamic-port:为true时 Jetty 会在portport + port-range之间选取第一个空闲端口;
  • port-range:动态端口搜索范围(默认100)。

注意hazelcast.yaml中的network.rest-api.enabled并不能替代上述 Jetty 开关。若配置了context-path: /seatunnel,所有 REST 端点都会移动该前缀之下(如/seatunnel/overview)。hazelcast.yaml的开关与 Jetty 相互独立,排查http://<host>:8080/不可达时,应首先确认上述enable-http/enable-https是否真正生效。

动态端口开启时实际端口可能不是 8080,请以启动日志中的SeaTunnel REST service will start on port xxx为准。下文所有示例均以http://<master>:8080为例,请替换为实际 master 地址与端口。

相关源码佐证

HTTP 相关配置项定义于 ServerConfigOptions.java(PORTENABLE_HTTPENABLE_DYNAMIC_PORTPORT_RANGECONTEXT_PATH等),全部 REST 端点路径常量集中定义于 RestConstant.java,Jetty 服务本身位于 JettyService.java。


2. 作业提交(Job Submission)

2.1 通过 JSON 请求体提交作业

curl -X POST http://<master>:8080/submit-job \ -H "Content-Type: application/json" \ -d @job.json

最小化的job.json结构(以 MySQL CDC → Console 为例):

{ "env": { "job.name": "my-cdc-job", "job.mode": "STREAMING", "checkpoint.interval": 30000 }, "source": [ { "plugin_name": "MySQL-CDC", "plugin_output": "mysql_cdc_result", "base-url": "jdbc:mysql://localhost:3306/mydb", "username": "cdc_user", "password": "password", "database-names": ["mydb"], "table-names": ["mydb.orders"], "startup.mode": "initial", "server-id": "5400-5404" } ], "transform": [], "sink": [ { "plugin_name": "Console", "plugin_input": ["mysql_cdc_result"] } ] }
  • env.job.mode:取BATCH(批式)或STREAMING(流式),checkpoint.interval单位为毫秒;
  • plugin_output/plugin_input:通过表名将上游 source 与下游 sink 串接起来,形成数据通路;
  • MySQL CDC 的server-id支持区间写法(如5400-5404),用于多并行度下分配唯一的 binlog 消费 ID。

2.2 带多个 Transform 的作业(JSON 格式)

{ "env": { "job.name": "etl-with-transforms", "job.mode": "BATCH" }, "source": [ { "plugin_name": "FakeSource", "plugin_output": "fake", "row.num": 100, "schema": { "fields": { "id": "int", "name": "string", "amount": "double" } } } ], "transform": [ { "plugin_name": "FieldMapper", "plugin_input": ["fake"], "plugin_output": "after_field_map", "field_mapper": { "id": "user_id", "name": "user_name" } }, { "plugin_name": "Filter", "plugin_input": ["after_field_map"], "plugin_output": "filtered", "fields": ["user_id", "user_name", "amount"] } ], "sink": [ { "plugin_name": "Console", "plugin_input": ["filtered"] } ] }

多个 transform 通过plugin_output/plugin_input依次串联成链:fake → FieldMapper → after_field_map → Filter → filtered → Console

2.3 提交响应

提交成功返回:

{ "jobId": "733584788375093248", "jobName": "my-cdc-job" }

请保存jobId,后续所有生命周期操作(查询、停止、恢复)均依赖它。

深入:请求体格式与底层提交链路

/submit-job的请求体除 JSON 外还支持 HOCON 与 SQL 两种格式,通过查询参数format指定(默认json)。从源码看,JobInfoService.submitJob 会依据ConfigFormat分别走ConfigFactory.parseString(HOCON)、SqlConfigBuilder.of(SQL)或RestUtil.buildConfig(JSON)三种解析路径,随后统一交给SeaTunnelServer完成作业提交。

此外还有两个提交相关端点:

  • POST /submit-job/upload:以上传配置文件方式提交(--form 'config_file=@"/temp/job.conf"'),支持.json(JSON)、.conf/.config(HOCON)、.sql(SQL 语法)三种文件;上传大小受seatunnel.engine.http.upload-max-file-size-mb(默认 10 MB)与upload-max-request-size-mb(默认 10 MB)限制,超限会在解析配置前直接拒绝,配置值 ≤ 0 表示不限。
  • POST /submit-jobs:批量提交,请求体为作业 JSON 数组,每个元素可携带params字段(内含jobIdjobNameisStartWithSavePoint)。

注意dryRun(试运行)功能刻意不在 REST API 中开放,仅在 SeaTunnel CLI 中可用;若通过 REST 传入dryRun参数,JobInfoService 会直接抛出IllegalArgumentException


3. 作业状态查询(Job Status Query)

3.1 查询单个作业详情

curl http://<master>:8080/job-info/<jobId>

响应字段:

字段说明
jobId唯一作业标识
jobName可读的作业名
jobStatusRUNNINGFINISHEDFAILEDCANCELLED
envOptions实际应用的 env 配置
createTime作业创建时间戳
jobDagDAG 结构(顶点与流水线边)
metricsSource/Sink 吞吐计数
finishedTime/errorMsg作业结束后返回
diagnostics作业运行时诊断信息(仅运行中、且可从 master 读取时返回)

字段返回规则(源码与文档双重确认)jobIdjobNamejobStatuscreateTimejobDagmetrics始终返回;envOptionspluginJarsUrlsisStartWithSavePoint仅在作业运行时返回;finishedTimeerrorMsg仅在作业结束后返回;diagnostics属于辅助信息,获取不到时该字段被省略而不会导致请求失败,且只有/job-info/:jobId返回它——/running-jobs不返回(为每个运行作业额外采集一次会多一次到 master 的往返)。

diagnostics中的pipelines[].restoreCount若在jobStatus保持RUNNING的同时持续增长,说明流水线处于崩溃重启循环(crash loop);maxRestoreCount对应job.retry.timesenv 选项设定的恢复上限。

3.2 查询所有运行中的作业

curl "http://<master>:8080/running-jobs?page=1&rows=10"

支持page(页码)与rows(每页条数)分页参数。

3.3 查询已结束作业

curl "http://<master>:8080/finished-jobs/FINISHED?page=1&rows=10"

state路径参数可取:FINISHEDFAILEDCANCELEDSAVEPOINT_DONEUNKNOWABLE。分页语义(见 REST API v2):提供page时响应包装为{"data": [...], "total": n};省略时返回裸数组;page/rows非正整数或越界会返回400

3.4 仅查询作业指标

curl http://<master>:8080/job-info/<jobId>

从响应metrics字段中读取关键指标:

指标含义
SourceReceivedCountSource 累计读取行数
SinkWriteCountSink 累计写入行数
SourceReceivedQPS当前读取吞吐(行/秒)
SinkWriteQPS当前写入吞吐(行/秒)

/job-info的 metrics 中还包含更完整的指标族:SourceReceivedBytes/SourceReceivedBytesPerSeconds(字节维度读写)、SinkCommittedCount/SinkCommittedQPS(checkpoint 成功后已提交行数与速率)、IntermediateQueueSize(算子间中间队列大小)、以及TableSourceReceived*TableSinkWrite*TableSinkCommitted*等按表(key 格式xxx#<table>)拆分的明细指标。这些指标名常量集中定义在 RestConstant.java。

补充:GET /running-job/:jobId为旧版端点,已被GET /job-info/:jobId取代并标记为 Deprecated,新代码请勿使用。


4. 查询作业日志(Querying Job Logs)

# 获取某个运行中作业日志的最后 N 行 curl "http://<master>:8080/logs/<jobId>"

该端点会跨所有节点汇总与指定jobId相关的日志:

  • GET /logs:返回全部节点的日志文件列表(默认 HTML 格式,?format=json可切换为 JSON);
  • GET /logs/<jobId>:跨所有节点检索指定作业的日志;
  • GET /logs/job-xxx.log:读取某个具体日志文件内容;
  • GET /log(单节点版本):从当前节点返回日志列表,http://localhost:5801/log即为 worker 节点上的日志入口。

对于日志文件分散在各自 worker 上的大规模部署,可直接使用 worker 自身的 REST 端口查询,或配置集中式日志(参见 Logging)。日志级别的运行时调整(/loggers端点)属于节点本地、重启即失效的临时覆盖,需要持久化的级别应写入config/log4j2.properties


5. 停止、取消与保存点语义(Stop, Cancel, and Savepoint)

三种操作的语义对比如下:

操作行为是否保留状态能否恢复
stop(优雅停止)等待在途数据冲刷完成在停止点做 checkpoint可以,通过--restore(REST 中为restoreMode
stop-with-savepoint优雅停止 + 显式写入保存点完整 savepoint可以,通过--restore
cancel(强制终止)立即终止不写入新状态仅能回到最近一次 checkpoint

5.1 优雅停止(不写保存点)

curl -X POST "http://<master>:8080/stop-job" \ -H "Content-Type: application/json" \ -d '{"jobId": "733584788375093248", "isStopWithSavePoint": false}'

5.2 带保存点停止

curl -X POST "http://<master>:8080/stop-job" \ -H "Content-Type: application/json" \ -d '{"jobId": "733584788375093248", "isStopWithSavePoint": true}'

保存点路径会打印在作业日志中,并出现在作业最终状态里,可用如下方式提取:

curl http://<master>:8080/job-info/733584788375093248 | \ python3 -c "import sys,json; d=json.load(sys.stdin); print(d.get('savepointPath', 'N/A'))"

5.3 强制取消(Cancel)

curl -X POST "http://<master>:8080/stop-job" \ -H "Content-Type: application/json" \ -d '{"jobId": "733584788375093248", "isStopWithSavePoint": false, "force": true}'

实现细节与注意事项

/stop-job/stop-jobs(批量停止,请求体为作业数组)分别由 StopJobServlet 与StopJobsServlet处理,最终委托给JobInfoService.stopJob/stopJobs(见 JobInfoService.java),成功后返回{"jobId": ...}

两点官方警告(记录于 REST API v2):

  • 若作业正处于DOING_SAVEPOINT状态且保存点未成功完成,使用force: true强制停止会将作业状态置为CANCELED
  • 强制停止可能遗留不完整、不一致的 checkpoint 数据,仅在异常/极端场景下使用

6. 作业恢复与重启(Job Recovery and Restart)

6.1 从最新 checkpoint 恢复

重新提交作业,并在查询参数中携带restoreMode=CHECKPOINTrestoreSourceJobId(指定要恢复的源作业 ID):

curl -X POST "http://<master>:8080/submit-job?restoreMode=CHECKPOINT&restoreSourceJobId=733584788375093248" \ -H "Content-Type: application/json" \ -d '{ "env": { "job.name": "my-cdc-job-restored", "job.mode": "STREAMING", "checkpoint.interval": 30000, "checkpoint.retain-after-job-cancelled": true }, "source": [ ... ], "sink": [ ... ] }'

同样的restoreModerestoreSourceJobId参数也适用于 上传配置文件提交端点——两个端点共享同一套恢复处理逻辑:

curl --location 'http://<master>:8080/submit-job/upload?restoreMode=CHECKPOINT&restoreSourceJobId=733584788375093248' \ --form 'config_file=@"/temp/my-cdc-job.conf"'

如果restoreSourceJobId对应的 checkpoint 数据缺失、已被清理或与当前作业不兼容,提交会快速失败(fail fast)

若希望被取消的作业仍可从 checkpoint 恢复,需要在取消之前以如下两种方式之一保留作业运行期间产生的 checkpoint 数据:

方式一:集群级默认配置(全局生效),写入config/seatunnel.yaml

seatunnel: engine: checkpoint: retain-after-job-cancelled: true

方式二:作业级 env 覆盖(仅当前作业生效),在 REST 请求体中配置:

{ "env": { "job.name": "my-cdc-job-restored", "job.mode": "STREAMING", "checkpoint.interval": 30000, "checkpoint.retain-after-job-cancelled": true }, "source": [ ... ], "sink": [ ... ] }

该选项默认值为false。若集群配置与作业 env 均未开启,被取消的作业默认仍会清理 checkpoint 数据;两者同时存在时,作业级 env 设置优先。该配置项在 ServerConfigOptions.java 中定义。

6.2 从最新保存点恢复

curl -X POST "http://<master>:8080/submit-job?restoreMode=SAVEPOINT&restoreSourceJobId=733584788375093248" \ -H "Content-Type: application/json" \ -d '{ "env": { "job.name": "my-cdc-job-restored", "job.mode": "STREAMING", "checkpoint.interval": 30000 }, "source": [ ... ], "sink": [ ... ] }'

6.3 从指定保存点路径恢复

curl -X POST http://<master>:8080/submit-job \ -H "Content-Type: application/json" \ -d '{ "env": { "job.name": "my-cdc-job-restored", "job.mode": "STREAMING", "checkpoint.interval": 30000, "restore.mode": "savepoint", "savepoint.path": "/seatunnel/checkpoint/savepoint/733584788375093248/1748595600000" }, "source": [ ... ], "sink": [ ... ] }'

深入:恢复参数如何被解析

从源码看,恢复逻辑集中在JobInfoServicevalidateCheckpointRestoreRequest中:restoreMode缺省时为RestoreMode.NONE,一旦restoreMode.isRestore()为真,则强制要求提供restoreSourceJobId,否则直接报错(见 JobInfoService.java)。此外,当只传isStartWithSavePoint而未指定restoreMode时,恢复源会回退到jobId参数(见 REST API v2 中/submit-job的参数说明)。关于 checkpoint/savepoint 的存储配置细节,可参见 Checkpoint Storage 与 State Storage and Recovery。


7. 认证与授权(Authentication and Authorization)

启用 Basic 认证后(配置方法见 Security),所有 REST API 调用都必须携带配置的用户名与密码:

curl -u admin:password "http://<master>:8080/running-jobs?page=1&rows=10"

未携带凭证时返回401 Unauthorized。HTTPS 的启用方式同样见 Security 文档。


8. REST API 性能考量(Performance Considerations)

8.1 大量已结束作业导致job-info变慢

finished-job-stateIMap 增长到数千条目规模时,/running-jobs/finished-jobs/:state端点会因全量扫描所有条目而变慢。缓解手段:

  1. 调小history-job-expire-minutes,缩短历史作业保留窗口(该配置项定义于 ServerConfigOptions.java,单位为分钟,同时驱动日志与历史状态的定期清理任务);
  2. 避免高频轮询 finished-jobs 端点,在监控层对结果做缓存;
  3. 监控看板直接按具体jobId查询,而不是列出全部作业。

8.2 并发提交速率

REST API 在 Hazelcast executor 线程池中同步处理提交。对于批量导入数百个作业的场景,应将提交速率控制在10~20 个/秒,避免压垮 master 节点。

8.3 动态端口分配

若开启enable-dynamic-port: true,不同 master 节点可能使用不同端口。可以从任意可达的 master 上调用 overview 端点查看集群实时状态:

# 从可达的 master 查看集群状态 curl http://<master>:8080/overview | \ python3 -c "import sys,json; print(json.load(sys.stdin))"

/overview返回集群层面的projectVersiontotalSlotunassignedSlotworksrunningJobspendingJobsfinishedJobsfailedJobscancelledJobs等摘要;若使用了动态 slot,totalSlotunassignedSlot恒为0。集群规模较大时,还可以结合/resource/workers查看各 worker 的 slot、CPU、内存与运行中作业快照。


9. 常见错误与故障排查(Common Errors and Troubleshooting)

错误原因解决办法
任何端点HTTP 404REST API 未启用或端口不对设置enable-http: true并核对端口
Connection refusedmaster 未启动或防火墙拦截端口确认 master 进程在运行;检查防火墙
job-info中提示jobId not found作业已结束或从未启动用预期的最终状态查询/finished-jobs/:state
提交返回400 Bad RequestJSON 格式错误或缺必填字段校验 JSON;检查plugin_name拼写
Job already exists with same job.id未先停止就重复提交相同的job.id先取消/停止已有作业,再重新提交
Unauthorized 401已开启 Basic 认证但未携带凭证请求中加上-u user:pass
Savepoint path not found保存点已被删除或路径错误检查 checkpoint 存储并给出正确路径

其他排查要点:

  • enable-dynamic-port生效,/overview或启动日志能帮助你确定真实端口,避免误判 404;
  • 处于DOING_SAVEPOINT的作业在保存点失败后如需强制终止,注意force: true会使作业进入CANCELED
  • 恢复提交失败多为 checkpoint/savepoint 数据缺失或不兼容,请先确认存储目录与restoreSourceJobId正确性。

See Also(延伸阅读)

  • REST API v2 Reference:全部端点的请求/响应结构、分页与参数完整参考
  • REST API v1 Reference:旧版 API 说明
  • Security Configuration:Basic 认证与 HTTPS 配置
  • Checkpoint Storage:checkpoint/savepoint 存储后端配置
  • State Storage and Recovery:状态存储与恢复机制详解
  • Logging:日志配置与集中式日志方案
  • CDC Pipeline Architecture:CDC 作业整体架构
  • 相关源码:JettyService.java、RestConstant.java、JobInfoService.java、StopJobServlet.java

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

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

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

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

立即咨询