资讯动态

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

发布时间:2026/9/17 21:30:57 来源:尧图企业网站定制
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 SeaTunnelZeta 引擎下作业的提交、状态查询、日志获取、停止/取消/保存点、恢复重启、认证与性能调优等完整生命周期操作。读者学完后可熟练使用 curl 与 Zeta 引擎内置 REST 服务完成一套可复制、可上线的作业管理流程并理解底层JobInfoService、Jetty 内嵌服务等实现原理。适用前提本 API 由 SeaTunnel EngineZeta内嵌的 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固定监听端口默认8080enable-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.javaPORT、ENABLE_HTTP、ENABLE_DYNAMIC_PORT、PORT_RANGE、CONTEXT_PATH等全部 REST 端点路径常量集中定义于 RestConstant.javaJetty 服务本身位于 JettyService.java。2. 作业提交Job Submission2.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.parseStringHOCON、SqlConfigBuilder.ofSQL或RestUtil.buildConfigJSON三种解析路径随后统一交给SeaTunnelServer完成作业提交。此外还有两个提交相关端点POST /submit-job/upload以上传配置文件方式提交--form config_file/temp/job.conf支持.jsonJSON、.conf/.configHOCON、.sqlSQL 语法三种文件上传大小受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 Query3.1 查询单个作业详情curl http://master:8080/job-info/jobId响应字段字段说明jobId唯一作业标识jobName可读的作业名jobStatusRUNNING、FINISHED、FAILED、CANCELLED等envOptions实际应用的 env 配置createTime作业创建时间戳jobDagDAG 结构顶点与流水线边metricsSource/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 loopmaxRestoreCount对应job.retry.timesenv 选项设定的恢复上限。3.2 查询所有运行中的作业curl http://master:8080/running-jobs?page1rows10支持page页码与rows每页条数分页参数。3.3 查询已结束作业curl http://master:8080/finished-jobs/FINISHED?page1rows10state路径参数可取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字段中读取关键指标指标含义SourceReceivedCountSource 累计读取行数SinkWriteCountSink 累计写入行数SourceReceivedQPS当前读取吞吐行/秒SinkWriteQPS当前写入吞吐行/秒/job-info的 metrics 中还包含更完整的指标族SourceReceivedBytes/SourceReceivedBytesPerSeconds字节维度读写、SinkCommittedCount/SinkCommittedQPScheckpoint 成功后已提交行数与速率、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 格式?formatjson可切换为 JSONGET /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可以通过--restoreREST 中为restoreModestop-with-savepoint优雅停止 显式写入保存点完整 savepoint可以通过--restorecancel强制终止立即终止不写入新状态仅能回到最近一次 checkpoint5.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; djson.load(sys.stdin); print(d.get(savepointPath, N/A))5.3 强制取消Cancelcurl -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 Restart6.1 从最新 checkpoint 恢复重新提交作业并在查询参数中携带restoreModeCHECKPOINT与restoreSourceJobId指定要恢复的源作业 IDcurl -X POST http://master:8080/submit-job?restoreModeCHECKPOINTrestoreSourceJobId733584788375093248 \ -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?restoreModeCHECKPOINTrestoreSourceJobId733584788375093248 \ --form config_file/temp/my-cdc-job.conf如果restoreSourceJobId对应的 checkpoint 数据缺失、已被清理或与当前作业不兼容提交会快速失败fail fast。若希望被取消的作业仍可从 checkpoint 恢复需要在取消之前以如下两种方式之一保留作业运行期间产生的 checkpoint 数据方式一集群级默认配置全局生效写入config/seatunnel.yamlseatunnel: 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?restoreModeSAVEPOINTrestoreSourceJobId733584788375093248 \ -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?page1rows10未携带凭证时返回401 Unauthorized。HTTPS 的启用方式同样见 Security 文档。8. REST API 性能考量Performance Considerations8.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 线程池中同步处理提交。对于批量导入数百个作业的场景应将提交速率控制在1020 个/秒避免压垮 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等摘要若使用了动态 slottotalSlot与unassignedSlot恒为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:passSavepoint 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 ConfigurationBasic 认证与 HTTPS 配置Checkpoint Storagecheckpoint/savepoint 存储后端配置State Storage and Recovery状态存储与恢复机制详解Logging日志配置与集中式日志方案CDC Pipeline ArchitectureCDC 作业整体架构相关源码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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价