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: 100enable-http:是否启动 HTTP 服务(代码默认false,打包的seatunnel.yaml示例已默认开启);port:固定监听端口(默认8080);enable-dynamic-port:为true时 Jetty 会在port到port + 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(PORT、ENABLE_HTTP、ENABLE_DYNAMIC_PORT、PORT_RANGE、CONTEXT_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字段(内含jobId、jobName、isStartWithSavePoint)。
注意:dryRun(试运行)功能刻意不在 REST API 中开放,仅在 SeaTunnel CLI 中可用;若通过 REST 传入dryRun参数,JobInfoService 会直接抛出IllegalArgumentException。
3. 作业状态查询(Job Status Query)
3.1 查询单个作业详情
curl http://<master>:8080/job-info/<jobId>响应字段:
| 字段 | 说明 |
|---|---|
jobId | 唯一作业标识 |
jobName | 可读的作业名 |
jobStatus | RUNNING、FINISHED、FAILED、CANCELLED等 |
envOptions | 实际应用的 env 配置 |
createTime | 作业创建时间戳 |
jobDag | DAG 结构(顶点与流水线边) |
metrics | Source/Sink 吞吐计数 |
finishedTime/errorMsg | 作业结束后返回 |
diagnostics | 作业运行时诊断信息(仅运行中、且可从 master 读取时返回) |
字段返回规则(源码与文档双重确认):jobId、jobName、jobStatus、createTime、jobDag、metrics始终返回;envOptions、pluginJarsUrls、isStartWithSavePoint仅在作业运行时返回;finishedTime、errorMsg仅在作业结束后返回;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路径参数可取:FINISHED、FAILED、CANCELED、SAVEPOINT_DONE、UNKNOWABLE。分页语义(见 REST API v2):提供page时响应包装为{"data": [...], "total": n};省略时返回裸数组;page/rows非正整数或越界会返回400。
3.4 仅查询作业指标
curl http://<master>:8080/job-info/<jobId>从响应metrics字段中读取关键指标:
| 指标 | 含义 |
|---|---|
SourceReceivedCount | Source 累计读取行数 |
SinkWriteCount | Sink 累计写入行数 |
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=CHECKPOINT与restoreSourceJobId(指定要恢复的源作业 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": [ ... ] }'同样的restoreMode与restoreSourceJobId参数也适用于 上传配置文件提交端点——两个端点共享同一套恢复处理逻辑:
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": [ ... ] }'深入:恢复参数如何被解析
从源码看,恢复逻辑集中在JobInfoService的validateCheckpointRestoreRequest中: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端点会因全量扫描所有条目而变慢。缓解手段:
- 调小
history-job-expire-minutes,缩短历史作业保留窗口(该配置项定义于 ServerConfigOptions.java,单位为分钟,同时驱动日志与历史状态的定期清理任务); - 避免高频轮询 finished-jobs 端点,在监控层对结果做缓存;
- 监控看板直接按具体
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返回集群层面的projectVersion、totalSlot、unassignedSlot、works、runningJobs、pendingJobs、finishedJobs、failedJobs、cancelledJobs等摘要;若使用了动态 slot,totalSlot与unassignedSlot恒为0。集群规模较大时,还可以结合/resource/workers查看各 worker 的 slot、CPU、内存与运行中作业快照。
9. 常见错误与故障排查(Common Errors and Troubleshooting)
| 错误 | 原因 | 解决办法 |
|---|---|---|
任何端点HTTP 404 | REST API 未启用或端口不对 | 设置enable-http: true并核对端口 |
Connection refused | master 未启动或防火墙拦截端口 | 确认 master 进程在运行;检查防火墙 |
job-info中提示jobId not found | 作业已结束或从未启动 | 用预期的最终状态查询/finished-jobs/:state |
提交返回400 Bad Request | JSON 格式错误或缺必填字段 | 校验 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),仅供参考