Flink REST API 完整指南:监控接口、异步操作与扩展机制
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
导读
Flink 内置了一套 REST-ful 风格的监控 API,用于查询正在运行作业以及最近完成作业的状态与统计信息。这套 API 既是 Flink Web 仪表盘(Dashboard)的数据来源,也开放给自定义监控工具直接调用。本文将以当前仓库中的 REST API 官方文档 为骨架,结合 flink-runtime 项目源码与仓库自带的 OpenAPI 规范文件,系统讲解 REST 服务端架构、端口与配置、版本化规则、异步操作(如 savepoint 触发)的正确用法,以及如何为 Flink 扩展新的 REST 请求端点。读完本文,你将能够熟练使用 curl 调用 JobManager 监控接口,理解 triggerid 异步模式的内部原理,并掌握添加自定义 REST handler 的完整开发步骤。
概览:监控 API 由谁提供,监听哪个端口
Flink 的监控 API 由作为JobManager一部分运行的 web 服务器提供支持。默认情况下,该服务器监听8081端口,端口号通过 Flink 配置文件 中的rest.port配置项进行修改。
从源码看,该配置项定义在 RestOptions.java 中,类型为整型、默认值 8081,并且继承了旧的web.port配置键作为弃用别名。需要注意它的两个细节:
- 文档描述明确提到:只有当高可用配置为 NONE 时,
rest.port才会被 REST 服务器真正用于绑定端口;如果启用了 HA,则优先使用rest.bind-port,因此生产环境启用 HA 时通常需要同时配置这两个键; rest.port同时是客户端(CLI、自定义工具)连接 JobManager REST 接口时使用的端口。
另一个关键事实是:监控 API 的 web 服务器和 Web 仪表盘的 web 服务器目前是同一个,它们在同一端口一起运行,但响应不同的 HTTP URL。也就是说,打开http://<jobmanager>:8081看到的是仪表盘页面,而向同一端口发送GET /v1/jobs/overview之类的请求则返回 JSON 监控数据。
此外,文档明确指出:在多个 JobManager(为了高可用)的情况下,每个 JobManager 都会运行自己的监控 API 实例,只有当某个 JobManager 被选举成为集群 leader 时,该实例才会提供已完成和正在运行作业的相关信息。这与WebMonitorEndpoint的 leader 选举语义一致——REST 服务端在成为 leader 后才对外提供有意义的集群视图。
与 REST 服务端相关的关键配置项
除rest.port外,围绕 REST 端点还有一组常见配置,均定义在 RestOptions.java 中,部分内容见 配置文件 的 "Advanced Options for the REST endpoint and Client" 一节。这里结合源码列出几个最常用项:
| 配置键 | 默认值 | 作用(依据源码注释与文档) |
|---|---|---|
rest.address | 空 | 客户端连接 Flink 使用的地址;应设置为 JobManager 所在主机名,或 Kubernetes 中位于 JobManager REST 接口前面的 Service 主机名 |
rest.port | 8081 | 客户端连接使用的端口;当未指定rest.bind-port时 REST 服务器绑定到该端口,且仅在高可用配置为 NONE 时生效 |
rest.bind-address | 空 | REST 服务器实际绑定的地址 |
rest.bind-port | 空 | REST 服务器实际绑定的端口,启用 HA 时优先于此值 |
rest.connection-timeout | 依赖默认实现 | REST 客户端连接超时时间 |
rest.retry.max-attempts | 20 | 可重试操作失败后客户端的最大重试次数(RestOptions.java) |
rest.retry.delay | 依赖默认实现 | 重试之间的延迟时间 |
rest.await-leader-timeout | 30 秒 | 客户端等待 leader 地址(如 Dispatcher 或 WebMonitorEndpoint)的最长时间(RestOptions.java) |
rest.async.store-duration | 5 分钟 | 异步操作结果在服务端保留的最长时间,超过后无法再查询(见下文"异步操作"章节) |
依据 RestOptions.java,
rest.flamegraph.stack-depth(默认 100)控制火焰图堆栈最大深度,rest.profiling.enabled(默认 false)控制实验性 profiler 功能开关,这些配置同样归属于 REST 专家选项区。
架构与扩展:REST 后端如何工作
REST API 的后端位于flink-runtime项目中,核心类是org.apache.flink.runtime.webmonitor.WebMonitorEndpoint,它负责配置服务器和请求路由。该类定义于 WebMonitorEndpoint.java,继承自org.apache.flink.runtime.rest.RestServerEndpoint。
从源码结构看,Flink 使用Netty和Netty Router库来处理 REST 请求和 URL 转换。官方文档给出的选型理由是该组合的依赖很轻量,且 Netty HTTP 性能非常好。RestServerEndpoint负责启动 Netty 服务器、绑定端口;WebMonitorEndpoint则在initializeHandlers()方法中完成全部请求路由的注册。
initializeHandlers()定义在 WebMonitorEndpoint.java,它返回一个List<Tuple2<RestHandlerSpecification, ChannelInboundHandler>>——即"请求规格 → 处理器"的映射列表。方法内部依次创建并注册了ClusterOverviewHandler(集群概览)、DashboardConfigHandler(仪表盘配置)、JobIdsHandler(作业 ID 列表)、JobStatusHandler(作业状态)、JobsOverviewHandler(作业概览)、ClusterConfigHandler(集群配置)等大量 handler。通过阅读该方法的导入列表,可以看到当前仓库中已经注册的 handler 家族,包括但不限于:
- 集群类:
/cluster(关闭集群)、/config(WebUI 配置)、/jobmanager/logs、/jobmanager/thread-dump等; - 作业类:
/jobs、/jobs/overview、/jobs/:jobid、/jobs/:jobid/status、/jobs/:jobid/exceptions、/jobs/:jobid/plan、/jobs/:jobid/execution-result等; - 检查点与 savepoint 类:
/jobs/:jobid/checkpoints/**、/jobs/:jobid/savepoints、/jobs/:jobid/rescaling等; - 指标类:
/jobmanager/metrics、/jobs/metrics、/jobs/:jobid/metrics、/taskmanagers/:taskmanagerid/metrics等。
如何添加一个新的 REST 请求
官方文档给出了扩展 REST API 的三步标准流程:
- 添加一个新的
MessageHeaders类,作为新请求的接口定义(URL、HTTP 方法、请求/响应类型); - 添加一个新的
AbstractRestHandler类,接收并处理该MessageHeaders描述的请求; - 将 handler 注册到
org.apache.flink.runtime.webmonitor.WebMonitorEndpoint#initializeHandlers()中。
文档推荐的范例是org.apache.flink.runtime.rest.handler.job.JobExceptionsHandler,它使用org.apache.flink.runtime.rest.messages.JobExceptionsHeaders。这两个类在仓库中分别位于:
- JobExceptionsHandler.java:继承自
AbstractExecutionGraphHandler,负责返回作业最近处理过的异常信息,并实现JsonArchivist接口(说明其结果会被历史归档机制保存); - JobExceptionsHeaders.java:其中定义了 URL 常量
"/jobs/:jobid/exceptions",getTargetRestEndpointURL()返回该 URL。
MessageHeaders中 URL 路径参数(如:jobid)会在请求匹配时被解析为路径参数对象(如JobIDPathParameter),handler 通过HandlerRequest获取这些参数后调用RestfulGateway查询对应的作业信息。这种"Headers 定义契约、Handler 实现逻辑、Endpoint 统一注册"的三层结构是理解 Flink REST 扩展机制的关键。
API 使用规则:版本化、默认版本与错误语义
REST API 是版本化的,通过在 URL 前面加上版本前缀来查询指定版本。前缀格式始终为v[version_number]。例如,要访问版本 1 的/foo/bar接口,需要请求:
/v1/foo/bar三条重要规则需要牢记:
- 未指定版本时,Flink 默认使用支持该请求的最旧版本——这意味着客户端如果不写前缀也能访问,但可能拿到的是旧版本语义;
- 查询不支持或不存在的版本会返回
404错误——不存在"静默降级",版本号写错会直接暴露出来; - 当前仓库的 OpenAPI 规范标题为
Flink JobManager REST API,版本标识为v1/2.0-SNAPSHOT(见 rest_v1_dispatcher.yml),即当前只提供 v1 版本。
异步操作与 triggerid 模式
这些 API 中存在多种异步操作,例如trigger savepoint、rescale a job。它们的行为模式是一致的:
- 客户端向触发类端点发起
POST请求; - 服务端立即返回一个
triggerid来标识这次 POST 操作; - 客户端随后使用该
triggerid轮询查询该操作的状态。
以 savepoint 为例(依据 rest_v1_dispatcher.yml):
POST /v1/jobs/:jobid/savepoints:触发一个 savepoint(可选的,之后是否取消作业),这是异步操作,返回202和TriggerResponse(内含triggerid);GET /v1/jobs/:jobid/savepoints/:triggerid:查询指定 savepoint 操作的进度与结果,返回200和AsynchronousOperationResult。
此外,对于stop-with-savepoint这类操作,你可以在触发请求的 body 中自行设置triggerId(英文版文档 docs/content/docs/ops/rest_api.md 对此有明确说明)。这意味着你可以安全地重试该操作而不会触发多次 savepoint——重试的请求携带同一个triggerid,服务端即可识别出这是同一操作。
但重试的安全性是有时限的:只有当rest.async.store-duration配置的异步操作存储时长尚未过期之前,重试才是安全的。该配置项定义于 RestOptions.java,默认值为5 分钟,含义是"异步操作结果在服务端存储的最大时长,一旦过期,操作结果将无法再被查询"。也就是说,如果你在 5 分钟之后用同一个triggerid重试,服务端可能已经丢失了原操作的记录,从而可能再次真正触发一个 savepoint。在设计自动化重试逻辑时,务必把超时窗口考虑进去。
JobManager API 参考与 OpenAPI 规范
当前仓库为 JobManager REST API 提供了OpenAPI 3.0.1 规范文件:docs/static/generated/rest_v1_dispatcher.yml。需要提醒的是,官方文档明确标注:OpenAPI specification 目前仍是实验性的(experimental),实际行为应以源码与运行时为准。
以该规范文件为索引,可以快速盘点 v1 版 JobManager 端点的主要分组(路径均基于规范文件paths段确认):
| 分组 | 代表端点 | 用途 |
|---|---|---|
| 集群 | DELETE /cluster | 关闭整个集群 |
| 配置 | GET /config | 返回 WebUI 的配置 |
| 数据集合 | GET /datasets、DELETE /datasets/{datasetid} | 查看与删除集群数据集(删除为异步操作) |
| Jar 管理 | GET /jars、POST /jars/upload、DELETE /jars/{jarid} | 列出、上传、删除用户 Jar |
| JobManager 信息 | GET /jobmanager/config、/jobmanager/environment、/jobmanager/logs、/jobmanager/metrics、/jobmanager/thread-dump | 查看 JM 配置、环境、日志、指标、线程转储 |
| 作业列表 | GET /jobs、GET /jobs/overview | 列出作业 ID 与作业概览 |
| 作业详情 | GET /jobs/{jobid}、/jobs/{jobid}/status、/jobs/{jobid}/config、/jobs/{jobid}/plan、/jobs/{jobid}/exceptions、/jobs/{jobid}/execution-result | 作业执行详情、状态、配置、执行计划、异常、执行结果 |
| 检查点 | GET /jobs/{jobid}/checkpoints、/jobs/{jobid}/checkpoints/config、/jobs/{jobid}/checkpoints/details/{checkpointid} | 检查点统计、配置与详情 |
| Savepoint / Rescale | POST /jobs/{jobid}/savepoints、GET /jobs/{jobid}/savepoints/{triggerid}、POST /jobs/{jobid}/rescaling、GET /jobs/{jobid}/rescaling/{triggerid} | 触发与查询异步操作 |
| 指标 | GET /jobmanager/metrics、/jobs/metrics、/jobs/{jobid}/metrics、/taskmanagers/{taskmanagerid}/metrics | 各类指标查询 |
| 顶点(算子) | GET /jobs/{jobid}/vertices/{vertexid}、/backpressure、/flamegraph、/subtasks/** | 顶点详情、背压、火焰图、子任务信息 |
| TaskManager | GET /taskmanagers、GET /taskmanagers/{taskmanagerid} | TM 列表与详情、日志、线程转储 |
其中/jobs/{jobid}/exceptions端点(对应上文示例 handler)支持两个查询参数:
maxExceptions:整型,限制返回异常条数的上限;failureLabelFilter:key:value形式的过滤集合,只返回带有全部指定 failure label 的异常。
规范中还注明了web.exception-history-size配置控制后端为每个作业收集的最近异常数量上限。
一个完整的 curl 调用示例
结合上述端点,给出一个可直接运行的调用序列(将<jobmanager>替换为实际主机名,作业 ID 可从/jobs列表获得):
# 1. 列出所有作业 ID(未指定版本,Flink 默认使用最旧支持版本) curl http://<jobmanager>:8081/jobs/overview # 2. 显式指定 v1 版本查询某个作业的当前状态 curl http://<jobmanager>:8081/v1/jobs/<jobid>/status # 3. 查询作业异常(限制最多返回 5 条) curl "http://<jobmanager>:8081/v1/jobs/<jobid>/exceptions?maxExceptions=5" # 4. 触发 savepoint(异步操作,返回 triggerid) curl -X POST -H "Content-Type: application/json" \ -d '{"target-directory": "file:///tmp/savepoints", "cancel-job": false}' \ http://<jobmanager>:8081/v1/jobs/<jobid>/savepoints # 5. 使用返回的 triggerid 查询 savepoint 操作状态 curl http://<jobmanager>:8081/v1/jobs/<jobid>/savepoints/<triggerid>结合源码理解端到端链路
把上面的知识点串起来,一次完整的监控查询请求在仓库内部的流转路径大致是:
- 请求到达 Netty HTTP 服务器(由
RestServerEndpoint启动),由 Netty Router 依据 URL 路由到WebMonitorEndpoint在initializeHandlers()中注册的 handler; MessageHeaders(如JobExceptionsHeaders)提供 URL 模板/jobs/:jobid/exceptions与请求/响应类型,负责将 HTTP 请求映射为类型化的HandlerRequest;AbstractRestHandler子类(如JobExceptionsHandler)从HandlerRequest中取出JobIDPathParameter等路径参数,通过GatewayRetriever获取当前 leader 的RestfulGateway,调用其方法从ExecutionGraph/ArchivedExecutionGraph中收集数据;- 结果封装为对应的
ResponseBody(如JobExceptionsInfoWithHistory)序列化为 JSON 返回。
对于异步操作(savepoint / rescaling),服务端将操作结果保存在一个受rest.async.store-duration控制的存储中,GET .../{triggerid}端点再从该存储中取出结果返回——这也是"存储时长过期后结果不可查"这一语义的来源。整个过程印证了官方文档对架构的概述:轻量依赖(Netty + Netty Router)、由WebMonitorEndpoint统一路由、MessageHeaders与AbstractRestHandler配对注册。
小结
Flink 的 REST 监控 API 是连接 Web 仪表盘、CLI 与自定义运维工具的统一数据面。理解其端口与绑定配置(rest.port/rest.bind-port在 HA 下的差异)、版本化规则(v1前缀、默认最旧版本、404 语义)、异步操作模式(triggerid + 轮询 + 存储时长窗口)是日常使用的核心;而MessageHeaders→AbstractRestHandler→initializeHandlers()三步注册流程,则为有定制需求的团队提供了清晰的扩展入口。当前仓库自带的 OpenAPI 规范 是快速检索全部 v1 端点最便捷的索引,建议配合 REST API 官方文档 与 配置文件说明 一起使用。
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考