ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Presto Statement Resource 完整指南:深入解析 /v1/statement 协议与查询执行链路

Presto Statement Resource 完整指南:深入解析 /v1/statement 协议与查询执行链路 大数据数据库后端【免费下载链接】prestoThe official home of the Presto distributed SQL query engine for big data项目地址https://gitcode.com/gh_mirrors/pre/presto点击查看免费下载导读Statement Resource 是 Presto 分布式 SQL 查询引擎面向客户端暴露的核心 HTTP 协议入口无论你使用 Presto CLI、JDBC 驱动还是自研客户端SQL 查询最终都会以 HTTP 请求的形式提交到协调节点Coordinator的/v1/statement端点并通过“提交—轮询—取结果”三步交互完成整条查询的生命周期管理。本文以官方文档 presto-docs/src/main/sphinx/rest/statement.rst 为主体结合当前仓库中 QueuedStatementResource.java、ExecutingStatementResource.java 与 StatementClientV1.java 的源码实现带你掌握 Statement 协议的全部端点、请求/响应格式、nextUri 轮询机制以及当前版本协议的最新演化细节。一、Statement Resource 在 Presto 架构中的位置Presto 采用经典的“协调节点Coordinator 工作节点Worker”架构。当你在终端里执行presto-cli --server localhost:8080 --catalog jmx --schema jmx并敲下一条 SQL 时CLI 并不是直接把 SQL 塞给某个执行引擎而是将它包装成一次 HTTP POST 请求发送到协调节点上的 Statement Resource协调节点负责解析 SQL、生成分布式执行计划、划分 Stage 与 Split工作节点负责实际读取数据、执行算子并把结果回传给协调节点Statement Resource 则负责在客户端与查询执行引擎之间搭建一座“轮询式”的桥梁。正如官方文档所言“The Presto client executes queries on behalf of a user against a catalog and a schema. When you run a query with the Presto CLI, it is calling out to the statement resource on the Presto coordinator.” 也就是说Statement Resource 是 CLI、JDBC 等所有官方客户端与 Presto 服务器交互的公共协议面。客户端侧的证据非常直接StatementClientV1.java 中buildQueryRequest方法把请求 URL 固定构造为/v1/statement并将 SQL 以文本形式 POST 出去HttpUrl url HttpUrl.get(session.getServer()); url url.newBuilder().encodedPath(/v1/statement).build(); Request.Builder builder prepareRequest(url) .post(RequestBody.create(MEDIA_TYPE_TEXT, query));由此可见Statement Resource 是理解整个 Presto 客户端协议的关键入口。二、协议总览四个端点、一条查询生命周期官方文档定义了 Statement Resource 的核心操作当前版本Path(/)前缀下在源码中对应如下端点方法路径源码位置作用POST/v1/statementQueuedStatementResource.java提交新语句Presto 生成 queryId 与 slugPUT/v1/statement/{queryId}?slug{slug}同文件#L303-L340提交语句但由调用方预生成 queryId/slugGET/v1/statement/queued/{queryId}/{token}同文件#L414-L458轮询排队状态等待查询被派发GET/v1/statement/executing/{queryId}/{token}ExecutingStatementResource.java轮询执行进度并拉取结果数据DELETE/v1/statement/queued/{queryId}/{token}QueuedStatementResource.java取消排队中的查询DELETE/v1/statement/executing/{queryId}/{token}ExecutingStatementResource.java取消执行中的查询GET/v1/statement/queued/retry/{queryId}QueuedStatementResource.java重新提交一条已失败的可重试查询当前版本新增一条查询的典型生命周期为客户端POST /v1/statement提交 SQL服务端立即返回 200响应中包含id、infoUri、nextUri、stats含 rootStage等字段但此时查询可能尚未真正开始执行——Presto 采用惰性执行lazy execution查询只有在客户端第一次轮询时才会被真正调度/派发这一点在 QueuedStatementResource.java 的注释中明确写明客户端反复GET nextUri所指地址直到nextUri为null表示查询结束中途可通过DELETE /v1/statement/.../{queryId}/{token}取消查询。2.1 惰性执行Lazy Execution与两段式轮询当前仓库的源码把“排队”与“执行”拆成了两个资源QueuedStatementResource负责提交与排队阶段的轮询ExecutingStatementResource负责执行阶段的轮询。客户端首先从POST响应里拿到形如/v1/statement/queued/{queryId}/{token}的nextUri排队阶段结束后再切换到/v1/statement/executing/{queryId}/{token}。这种设计的核心目的是让协调节点不必为每个提交的查询立即占用调度线程而是把开销推迟到客户端真正索要结果时从而提升高并发下的吞吐能力。从源码看排队阶段轮询时服务端会等待“查询被派发”这一状态变化// wait for query to be dispatched, up to the wait timeout ListenableFuture? futureStateChange addTimeout( waitForDispatchedAsync, () - null, WAIT_ORDERING.min(MAX_WAIT_TIME, maxWait), timeoutExecutor);其中maxWait是客户端可选的maxWait查询参数服务端用“最长等待时间”进行长轮询long polling既减少空轮询对服务端的压力也降低客户端延迟。三、POST /v1/statement提交查询3.1 请求头X-Presto-* 系列官方文档列出的标准请求头如下它们在客户端侧由 PrestoHeaders.java 统一定义请求头说明是否必选X-Presto-User以哪个用户身份执行语句可选但生产环境强烈建议设置配合认证/授权X-Presto-Source查询来源标识如presto-cli用于资源组与审计可选X-Presto-Catalog查询目标 Catalog如jmx、tpch、hive可选X-Presto-Schema查询目标 Schema可选除上述之外PrestoHeaders.java 中还定义了X-Presto-Session会话属性形如keyvalue、X-Presto-Transaction-Id事务 ID、X-Presto-Prefix-Url代理前缀、X-Presto-Retry-Query标记重试查询等扩展头。Catalog/Schema/User 等信息会被服务端封装为HttpRequestSessionContext见QueuedStatementResource.postStatement中对new HttpRequestSessionContext(...)的调用构成查询执行的会话上下文。3.2 请求体与官方示例请求体就是一条纯文本 SQLContent-Type为text/plain。文档给出的完整示例POST /v1/statement HTTP/1.1 Host: localhost:8001 X-Presto-Catalog: jmx X-Presto-Source: presto-cli X-Presto-Schema: jmx User-Agent: StatementClient/0.55-SNAPSHOT X-Presto-User: tobrie1 Content-Length: 41 select name from java.lang:typeruntime注意请求体中的 SQL 不做任何转义包装直接以原始文本发送示例中的java.lang:typeruntime是 JMX Catalog 下的 MBean 对象名因此需要双引号包裹。3.3 服务端处理逻辑源码级postStatement的源码处理流程QueuedStatementResource.java大致为校验SQL 为空则返回 400SQL statement is emptyX-Presto-Prefix-Url非法则返回 400构建会话上下文HttpRequestSessionContext从 HTTP 请求头提取 User/Catalog/Schema/Session 属性等生成 queryId 与 slugdispatchManager.createQueryId()生成全局唯一的查询 IDcreateSlug()生成一个随机 nonce形如x 32 位小写十六进制用于后续轮询请求的身份校验防止查询结果被未授权方读取private static String createSlug() { return x randomUUID().toString().toLowerCase(ENGLISH).replace(-, ); }登记查询将Query对象放入内存中的queriesMap返回初始响应携带Cache-Control: max-age60Query.CACHE_CONTROL_MAX_AGE_SEC 60JSON 序列化QueryResults。3.4 响应体逐字段解析文档给出的响应示例已按当前版本字段对齐{ id: 20140108_110629_00011_dk5x2, infoUri: http://localhost:8001/v1/query/20140108_110629_00011_dk5x2, partialCancelUri: http://10.193.207.128:8080/v1/stage/20140108_110629_00011_dk5x2.1, nextUri: http://localhost:8001/v1/statement/20140108_110629_00011_dk5x2/1, columns: [ { name: name, type: varchar } ], stats: { state: RUNNING, scheduled: false, nodes: 1, totalSplits: 0, queuedSplits: 0, runningSplits: 0, completedSplits: 0, cpuTimeMillis: 0, wallTimeMillis: 0, processedRows: 0, processedBytes: 0, rootStage: { stageId: 0, state: SCHEDULED, done: false, nodes: 1, totalSplits: 0, queuedSplits: 0, runningSplits: 0, completedSplits: 0, cpuTimeMillis: 0, wallTimeMillis: 0, processedRows: 0, processedBytes: 0, subStages: [ { stageId: 1, state: SCHEDULED, done: false, nodes: 1, totalSplits: 0, queuedSplits: 0, runningSplits: 0, completedSplits: 0, cpuTimeMillis: 0, wallTimeMillis: 0, processedRows: 0, processedBytes: 0, subStages: [] } ] } } }关键字段含义字段含义id查询标识符格式为yyyyMMdd_HHmmss_xxxxx_yyyyy后续所有操作都要引用它infoUri查询详情地址指向/v1/query/{queryId}可用于获取查询的完整元数据信息partialCancelUri部分结果取消地址指向某个具体 StagenextUri最重要的字段客户端下一轮轮询的地址为null时表示查询结束columns结果列定义名称与类型stats查询执行统计核心是Stage 树stats.rootStage根 StagestageId 固定为0聚合所有工作节点上子 Stage 的结果并经协调节点交付给客户端关于 Stage 树官方文档特别强调“Every query has a root stage and the root stage is given a stage identifier of ‘0’”。根 Stage 是分布式执行计划的最顶层聚合节点其subStages数组递归描述了整棵执行树——例如上例中根 Stage0下挂着一个子 Stage1。从源码结构看当前版本中 Stage 树的统计由StatementStats与其内部的StageStats递归结构承载见 QueryResourceUtil.java 中StatementStats.builder()的构造逻辑。3.5 为什么是“chunked”响应示例响应头中的Transfer-Encoding: chunked是 HTTP/1.1 的分块传输编码由于查询结果可能持续产出服务端无法预先计算Content-Length因此采用分块方式边算边写。客户端侧的 StatementClientV1.java 用JsonResponse.execute(...)接收这个流式 JSON并做严格的状态检查只要状态码不是 200 或响应无值就判定CLIENT_ERROR并抛出requestFailedException。四、PUT /v1/statement/{queryId}?slug{slug}预生成 ID 的提交官方文档描述 PUT 是 POST 的“同构物”唯一区别是 queryId 与 slug 由调用方显式提供而非 Presto 自动生成。这一设计最初是为了支持幂等提交与客户端驱动的重试当网络抖动导致提交请求丢失时客户端可以用同一个 queryId 与 slug 重新 PUT服务端通过queries.computeIfAbsent(queryId, ...)保证同一 queryId 只登记一次。源码中的关键逻辑QueuedStatementResource.javaQuery attemptedQuery new Query(statement, sessionContext, dispatchManager, executingQueryResponseProvider, 0, queryId, slug); Query query queries.computeIfAbsent(queryId, unused - attemptedQuery); if (attemptedQuery ! query !attemptedQuery.getSlug().equals(query.getSlug()) || query.getLastToken() ! 0) { throw badRequest(CONFLICT, Query already exists); }即如果 queryId 已存在且 slug 不一致、或该查询已被轮询过lastToken ! 0服务端会返回 409 Conflict防止查询被重复执行或劫持。参数说明queryId路径参数预生成的查询标识符slug查询参数与该查询绑定的 nonce后续所有轮询/取消请求都必须携带它其余X-Presto-*请求头与 POST 完全一致。五、GET /v1/statement/.../{queryId}/{token}轮询与取数5.1 官方示例提交后客户端拿到nextUri然后不断 GET 它。文档示例GET /v1/statement/20140108_110629_00011_dk5x2/1 HTTP/1.1 Host: localhost:8001 User-Agent: StatementClient/0.55-SNAPSHOT注意示例 URL 是协议早期版本的路径格式在当前版本中nextUri会被构造成/v1/statement/queued/{queryId}/{token}或/v1/statement/executing/{queryId}/{token}见 QueryResourceUtil.getQueuedUri该函数把路径拼为/v1/statement/queued并追加queryId、token与slug参数。token 是一个单调递增的序号来自上一次响应的nextUri用于标识“下一批结果”slug 则用于鉴权。5.2 最终响应结果数据与完成态统计当查询完成时响应中会出现data数组并携带完成态统计。文档示例{ id: 20140108_110629_00011_dk5x2, infoUri: http://localhost:8001/v1/query/20140108_110629_00011_dk5x2, columns: [ { name: name, type: varchar } ], data: [ [4165domU-12-31-39-0F-CC-72] ], stats: { state: FINISHED, scheduled: true, nodes: 1, totalSplits: 2, queuedSplits: 0, runningSplits: 0, completedSplits: 2, cpuTimeMillis: 1, wallTimeMillis: 4, processedRows: 1, processedBytes: 27, rootStage: { stageId: 1, state: FINISHED, done: true, nodes: 1, totalSplits: 1, queuedSplits: 0, runningSplits: 0, completedSplits: 1, cpuTimeMillis: 0, wallTimeMillis: 0, processedRows: 1, processedBytes: 32, subStages: [ { stageId: 1, state: FINISHED, done: true, nodes: 1, totalSplits: 1, queuedSplits: 0, runningSplits: 0, completedSplits: 1, cpuTimeMillis: 0, wallTimeMillis: 4, processedRows: 1, processedBytes: 27, subStages: [] } ] } } }对比初始响应可以清晰看到生命周期演进stats.state由RUNNING变为FINISHEDscheduled由false变为truecompletedSplits追上totalSplitsdata数组携带真实结果行每行是一个与columns顺序对应的值数组根 Stage 的done变为true。这个示例也解释了为何nextUri机制是必须的结果是一批批返回的分块、分 token客户端必须持续轮询直到nextUri null。5.3 客户端侧如何驱动轮询源码级StatementClientV1.advance() 是客户端轮询循环的核心URI nextUri currentStatusInfo().getNextUri(); if (nextUri null) { state.compareAndSet(State.RUNNING, State.FINISHED); return false; } validateNextUriSource(nextUri, currentStatusInfo().getInfoUri()); Request request prepareRequest(HttpUrl.get(nextUri)).build();几个值得注意的细节终止条件nextUri null即查询结束无论成功还是失败客户端状态置为FINISHED安全校验validateNextUriSource会校验nextUri的 host 与 port 和infoUri一致防止被篡改的重定向地址劫持防 SSRF失败重试与退避轮询失败后客户端会按attempts * 100ms退避重试直到超过clientRequestTimeout才抛出异常服务端侧每次轮询都会做QueryRateLimiter 限流并以maxWait长轮询与targetResultSize目标结果大小上限为MAX_TARGET_RESULT_SIZE等参数控制响应节奏见 ExecutingStatementResource.getQueryResults。六、DELETE /v1/statement/.../{queryId}/{token}取消查询当客户端不再需要查询结果例如用户按下 CtrlC时应调用 DELETE 主动取消避免查询在集群上继续空转消耗资源。文档定义的语义与源码实现一致无论排队阶段/v1/statement/queued/{queryId}/{token}还是执行阶段/v1/statement/executing/{queryId}/{token}DELETE 都要求queryId、token以及鉴权用slug服务端找到对应查询后调用cancel()并返回204 No ContentDELETE Path(/v1/statement/queued/{queryId}/{token}) public Response cancelQuery( PathParam(queryId) QueryId queryId, PathParam(token) long token, QueryParam(slug) String slug) { getQuery(queryId, slug).cancel(); return Response.noContent().build(); }若 queryId 不存在或 slug 不匹配服务端返回404 Query not found。从源码结构看取消动作会逐级传播到DispatchManager与底层查询执行器从而终止各工作节点上正在运行的 task。七、协议的安全性设计Statement Resource 协议中内建了多层防护理解它们有助于自研客户端时避免踩坑slug 随机 nonce每个查询都绑定一个随机生成的 slugx UUID 去连字符共 33 字符。所有 GET/DELETE 轮询与取消请求都必须携带与 queryId 匹配的 slug服务端在 getQuery 中校验“查询存在且 slug 一致”否则返回 404。这有效防止了他人通过猜测 queryId 窃取查询结果token 单调递增token 标识结果批次序号服务端据此判断客户端请求是否合法推进重复 token 会被拒绝或重定向并配合lastToken检测 PUT 冲突X-Presto-User 认证服务端通过HttpRequestSessionContext提取用户身份配合 RolesAllowed(USER) 注解进行 RBAC 校验——Path(/)RolesAllowed(USER)意味着整个 Statement 资源只允许已认证的用户角色访问nextUri 源校验客户端侧如前所述validateNextUriSource 会核对 host 与 port防止恶意服务端把客户端导向外部地址X-Presto-Prefix-Url 校验服务端侧abortIfPrefixUrlInvalid会在构造响应前校验前缀 URL 的合法性非法则返回 400。八、从协议到客户端一次完整查询的旅程综合以上内容一次标准查询的完整时序如下客户端POST /v1/statement携带X-Presto-User/Catalog/Schema/Source头与 SQL 文本协调节点生成queryId、slug登记Query立即返回初始QueryResults含id、infoUri、nextUri、stats.rootStage客户端循环执行advance()若nextUri null结束否则 GETnextUri携带 slug排队阶段轮询QueuedStatementResource执行阶段轮询ExecutingStatementResource每个响应可能携带新的columns、data批次与最新的stats客户端收集所有data批次直到查询状态FINISHED如遇中断客户端DELETE取消服务端返回 204。这套“提交即返回、轮询取结果”的协议设计使得 Presto 能在大规模并行执行的同时保持客户端与服务端之间极简的 HTTP 交互也为后续的查询重试X-Presto-Retry-Query头与/v1/statement/queued/retry/{queryId}端点等高级特性预留了扩展空间。九、参考资料与延伸阅读官方 REST 文档statement.rst本文主体来源服务端实现QueuedStatementResource.java、ExecutingStatementResource.java、QueryResourceUtil.java客户端实现StatementClientV1.java、StatementClient.java、PrestoHeaders.java测试验证TestServer.java含/v1/statement相关端到端测试说明本文中nextUri的路径格式同时出现了文档中的早期格式/v1/statement/{queryId}/{token}与当前版本源码中的新格式/v1/statement/queued|executing/{queryId}/{token}。两者差异源于协议演进客户端不应硬编码路径而应始终以服务端返回的nextUri为准。赞分享大数据数据库后端【免费下载链接】prestoThe official home of the Presto distributed SQL query engine for big data项目地址https://gitcode.com/gh_mirrors/pre/presto点击查看免费下载相关推荐Presto Client REST API 完整指南从 /v1/statement 协议到跨集群查询重试Presto Client REST API 完整指南从 /v1/statement 协议到跨集群查询重试 本篇技术指南基于 Presto 官方文档与客户端/大数据数据库后端Presto REST API 完整指南/v1/node、/v1/query、/v1/stage、/v1/statement 与 /v1/task 端点详解Presto REST API 完整指南/v1/node、/v1/query、/v1/stage、/v1/statement 与 /v1/task 端点详解大数据数据库后端Presto Query Resource REST 接口全解析用 /v1/query 监控与诊断查询Presto Query Resource REST 接口全解析用 /v1/query 监控与诊断查询 导读 本文基于 Presto 官方文档中的 Query大数据数据库后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表