尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

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

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-address空REST 服务器实际绑定的地址rest.bind-port空REST 服务器实际绑定的端口启用 HA 时优先于此值rest.connection-timeout依赖默认实现REST 客户端连接超时时间rest.retry.max-attempts20可重试操作失败后客户端的最大重试次数RestOptions.javarest.retry.delay依赖默认实现重试之间的延迟时间rest.await-leader-timeout30 秒客户端等待 leader 地址如 Dispatcher 或 WebMonitorEndpoint的最长时间RestOptions.javarest.async.store-duration5 分钟异步操作结果在服务端保留的最长时间超过后无法再查询见下文异步操作章节依据 RestOptions.javarest.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它返回一个ListTuple2RestHandlerSpecification, ChannelInboundHandler——即请求规格 → 处理器的映射列表。方法内部依次创建并注册了ClusterOverviewHandler集群概览、DashboardConfigHandler仪表盘配置、JobIdsHandler作业 ID 列表、JobStatusHandler作业状态、JobsOverviewHandler作业概览、ClusterConfigHandler集群配置等大量 handler。通过阅读该方法的导入列表可以看到当前仓库中已经注册的 handler 家族包括但不限于集群类/cluster关闭集群、/configWebUI 配置、/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/exceptionsgetTargetRestEndpointURL()返回该 URL。MessageHeaders中 URL 路径参数如:jobid会在请求匹配时被解析为路径参数对象如JobIDPathParameterhandler 通过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.ymlPOST /v1/jobs/:jobid/savepoints触发一个 savepoint可选的之后是否取消作业这是异步操作返回202和TriggerResponse内含triggeridGET /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}列出、上传、删除用户 JarJobManager 信息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 / RescalePOST /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/**顶点详情、背压、火焰图、子任务信息TaskManagerGET /taskmanagers、GET /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?maxExceptions5 # 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()中注册的 handlerMessageHeaders如JobExceptionsHeaders提供 URL 模板/jobs/:jobid/exceptions与请求/响应类型负责将 HTTP 请求映射为类型化的HandlerRequestAbstractRestHandler子类如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),仅供参考
返回列表