news 2026/9/23 7:32:54

Flink REST API 完整指南:监控接口、异步操作与扩展机制

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink REST API 完整指南:监控接口、异步操作与扩展机制

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.port8081客户端连接使用的端口;当未指定rest.bind-port时 REST 服务器绑定到该端口,且仅在高可用配置为 NONE 时生效
rest.bind-addressREST 服务器实际绑定的地址
rest.bind-portREST 服务器实际绑定的端口,启用 HA 时优先于此值
rest.connection-timeout依赖默认实现REST 客户端连接超时时间
rest.retry.max-attempts20可重试操作失败后客户端的最大重试次数(RestOptions.java)
rest.retry.delay依赖默认实现重试之间的延迟时间
rest.await-leader-timeout30 秒客户端等待 leader 地址(如 Dispatcher 或 WebMonitorEndpoint)的最长时间(RestOptions.java)
rest.async.store-duration5 分钟异步操作结果在服务端保留的最长时间,超过后无法再查询(见下文"异步操作"章节)

依据 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 使用NettyNetty 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 的三步标准流程:

  1. 添加一个新的MessageHeaders,作为新请求的接口定义(URL、HTTP 方法、请求/响应类型);
  2. 添加一个新的AbstractRestHandler,接收并处理该MessageHeaders描述的请求;
  3. 将 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

三条重要规则需要牢记:

  1. 未指定版本时,Flink 默认使用支持该请求的最旧版本——这意味着客户端如果不写前缀也能访问,但可能拿到的是旧版本语义;
  2. 查询不支持或不存在的版本会返回404错误——不存在"静默降级",版本号写错会直接暴露出来;
  3. 当前仓库的 OpenAPI 规范标题为Flink JobManager REST API,版本标识为v1/2.0-SNAPSHOT(见 rest_v1_dispatcher.yml),即当前只提供 v1 版本。

异步操作与 triggerid 模式

这些 API 中存在多种异步操作,例如trigger savepointrescale a job。它们的行为模式是一致的:

  1. 客户端向触发类端点发起POST请求;
  2. 服务端立即返回一个triggerid来标识这次 POST 操作;
  3. 客户端随后使用该triggerid轮询查询该操作的状态。

以 savepoint 为例(依据 rest_v1_dispatcher.yml):

  • POST /v1/jobs/:jobid/savepoints:触发一个 savepoint(可选的,之后是否取消作业),这是异步操作,返回202TriggerResponse(内含triggerid);
  • GET /v1/jobs/:jobid/savepoints/:triggerid:查询指定 savepoint 操作的进度与结果,返回200AsynchronousOperationResult

此外,对于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 /datasetsDELETE /datasets/{datasetid}查看与删除集群数据集(删除为异步操作)
Jar 管理GET /jarsPOST /jars/uploadDELETE /jars/{jarid}列出、上传、删除用户 Jar
JobManager 信息GET /jobmanager/config/jobmanager/environment/jobmanager/logs/jobmanager/metrics/jobmanager/thread-dump查看 JM 配置、环境、日志、指标、线程转储
作业列表GET /jobsGET /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 / RescalePOST /jobs/{jobid}/savepointsGET /jobs/{jobid}/savepoints/{triggerid}POST /jobs/{jobid}/rescalingGET /jobs/{jobid}/rescaling/{triggerid}触发与查询异步操作
指标GET /jobmanager/metrics/jobs/metrics/jobs/{jobid}/metrics/taskmanagers/{taskmanagerid}/metrics各类指标查询
顶点(算子)GET /jobs/{jobid}/vertices/{vertexid}/backpressure/flamegraph/subtasks/**顶点详情、背压、火焰图、子任务信息
TaskManagerGET /taskmanagersGET /taskmanagers/{taskmanagerid}TM 列表与详情、日志、线程转储

其中/jobs/{jobid}/exceptions端点(对应上文示例 handler)支持两个查询参数:

  • maxExceptions:整型,限制返回异常条数的上限;
  • failureLabelFilterkey: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>

结合源码理解端到端链路

把上面的知识点串起来,一次完整的监控查询请求在仓库内部的流转路径大致是:

  1. 请求到达 Netty HTTP 服务器(由RestServerEndpoint启动),由 Netty Router 依据 URL 路由到WebMonitorEndpointinitializeHandlers()中注册的 handler;
  2. MessageHeaders(如JobExceptionsHeaders)提供 URL 模板/jobs/:jobid/exceptions与请求/响应类型,负责将 HTTP 请求映射为类型化的HandlerRequest
  3. AbstractRestHandler子类(如JobExceptionsHandler)从HandlerRequest中取出JobIDPathParameter等路径参数,通过GatewayRetriever获取当前 leader 的RestfulGateway,调用其方法从ExecutionGraph/ArchivedExecutionGraph中收集数据;
  4. 结果封装为对应的ResponseBody(如JobExceptionsInfoWithHistory)序列化为 JSON 返回。

对于异步操作(savepoint / rescaling),服务端将操作结果保存在一个受rest.async.store-duration控制的存储中,GET .../{triggerid}端点再从该存储中取出结果返回——这也是"存储时长过期后结果不可查"这一语义的来源。整个过程印证了官方文档对架构的概述:轻量依赖(Netty + Netty Router)、由WebMonitorEndpoint统一路由、MessageHeadersAbstractRestHandler配对注册。

小结

Flink 的 REST 监控 API 是连接 Web 仪表盘、CLI 与自定义运维工具的统一数据面。理解其端口与绑定配置(rest.port/rest.bind-port在 HA 下的差异)、版本化规则(v1前缀、默认最旧版本、404 语义)、异步操作模式(triggerid + 轮询 + 存储时长窗口)是日常使用的核心;而MessageHeadersAbstractRestHandlerinitializeHandlers()三步注册流程,则为有定制需求的团队提供了清晰的扩展入口。当前仓库自带的 OpenAPI 规范 是快速检索全部 v1 端点最便捷的索引,建议配合 REST API 官方文档 与 配置文件说明 一起使用。

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 7:32:51

3步搞定razer驱动:从报错到实战项目避坑指南

3步搞定razer驱动:从报错到实战项目避坑指南 报错堆成山,StackTrace 根本看不懂?别慌,这不仅是你的问题,更是很多开发者在接入硬件外设时的通病。当你在做一个 实战项目 ,需要调用雷蛇(Razer)键盘、鼠标或耳麦的高级功能时, razer驱动 的底层逻辑往往成了拦路虎。…

作者头像 李华
网站建设 2026/9/23 7:32:45

3个核心避坑指南搞定喷墨打印机连供逻辑

3个核心避坑指南搞定喷墨打印机连供逻辑 别再对着教程发呆,代码跑不通才是真痛点。很多老哥在掘金技术社区问连供系统,答案往往不在纸上,而在数据流里。 概念速懂 喷墨打印机连供不是换个墨盒那么简单,它是把外置墨仓通过细管连接到喷头,实现持续供墨。传统墨盒是“一次性电池”,连供是“充电宝”。…

作者头像 李华
网站建设 2026/9/23 7:32:34

3步搞定看火山小视频底层逻辑保姆级教程

3步搞定看火山小视频底层逻辑保姆级教程 别再去翻那厚达几百页的官方文档了,真的会劝退人。 官方文档太长抓不住重点,是很多转行开发的朋友最大的噩梦。 这篇 保姆级教程 ,直接带你拆解【看火山小视频】背后的技术真相,用3步讲透原理。…

作者头像 李华
网站建设 2026/9/23 7:32:31

清华紫光输入法手写实现:保姆级教程解决StackTrace报错

清华紫光输入法手写实现:保姆级教程解决StackTrace报错 面对满屏红色的 StackTrace,是不是感觉脑瓜子嗡嗡的?别慌,这种报错堆栈看着吓人,其实核心就卡在几个关键节点。今天这篇清华紫光输入法手写实现的保姆级教程,就是专门为你准备的。我们不讲虚的,直接拆解代码逻辑,带你从零搭建一个能跑、…

作者头像 李华
网站建设 2026/9/23 7:32:28

3个实战项目看透www.33qqbb.com原理面试不挂

3个实战项目看透www.33qqbb.com原理面试不挂 面试被问“www.33qqbb.com”底层原理,你张嘴卡壳?别慌,这不是你记忆力差,而是没人教你怎么把代码和原理对应起来。我在三个真实实战项目中踩过坑,发现只要抓住核心链路,这种问题根本难不倒你。 一句话原理:数据从哪来,到哪去…

作者头像 李华
网站建设 2026/9/23 7:32:21

别被忽悠!虾青素的作用图解原理与工程避坑指南

别被忽悠!虾青素的作用图解原理与工程避坑指南 面试被问“虾青素的作用”却支支吾吾答不上来?这在化工、食品甚至生物医药行业的招聘中太常见了。很多候选人只背了“抗氧化”三个字,面对追问就露馅。 今天咱们不整虚的,直接上 图解原理…

作者头像 李华