ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Flink SQL Gateway 完全指南:架构原理、启动配置与 REST 查询实战

Apache Flink SQL Gateway 完全指南:架构原理、启动配置与 REST 查询实战 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载SQL Gateway 是 Apache Flink 提供的一项服务它允许多个远程客户端并发地提交 SQL、执行查询、检索元数据并进行在线数据分析为 Flink Table 程序提供统一的服务化入口。本文将以 SQL Gateway 官方 Overview 文档为主体结合仓库源码flink-table/flink-sql-gateway模块深入讲解其可插拔架构、开箱即用的启动方式、REST 端点下的完整查询流程以及全部核心配置项帮助读者掌握从零搭建 SQL Gateway 服务并完成第一次远程 SQL 查询的完整链路。SQL Gateway 是什么从定位上看SQL Gateway 是一个独立的常驻服务它解决了传统在本地开发环境中逐个执行 SQL的低效模式让多个客户端能够从远程以并发方式提交 SQL。根据官方文档的定义它提供了一种便捷的方式来提交 Flink Job把 SQL 语句转换为可执行的 Flink 作业提交到集群查询元数据获取 Catalog、Database、Table、函数等元数据信息在线分析数据通过交互式查询实时分析数据。在整体设计上SQL Gateway 由可插拔的 Endpoint端点与SqlGatewayService服务处理器两部分构成。SqlGatewayService是所有 Endpoint 复用的核心处理器负责统一处理各类请求而 Endpoint 则是用户接入服务的入口用户可以根据 Endpoint 的类型选择不同的客户端工具来连接。这一设计在源码中有着清晰的印证。SQL Gateway 的进程入口 SqlGateway.java 在start()方法中完成了服务 端点的组合装配先启动SessionManager再创建SqlGatewayServiceImpl最后通过SqlGatewayEndpointFactoryUtils.createSqlGatewayEndpoint(sqlGatewayService, defaultConfig)根据配置实例化出对应的 Endpoint 并逐个启动public void start() throws Exception { sessionManager.start(); SqlGatewayService sqlGatewayService new SqlGatewayServiceImpl(sessionManager); try { endpoints.addAll( SqlGatewayEndpointFactoryUtils.createSqlGatewayEndpoint( sqlGatewayService, defaultConfig)); for (SqlGatewayEndpoint endpoint : endpoints) { endpoint.start(); } } catch (Throwable t) { // ... } }而服务端的核心能力面由 SqlGatewayService.java 这个PublicEvolving接口定义它按功能划分了以下几组 API分组能力Session 管理openSession/closeSession/configureSession/getSessionConfig/getSessionEndpointVersionOperation 管理submitOperation/cancelOperation/closeOperation/getOperationInfo/getOperationResultSchemaStatement 执行executeStatement/fetchResults支持 token 与FetchOrientation两种拉取方式Catalog 查询getCurrentCatalog/listCatalogs/listDatabases/listTables/getTable函数查询listUserDefinedFunctions/listSystemFunctions/getFunctionDefinition工具能力getGatewayInfo/completeStatementSQL 自动补全提示物化表refreshMaterializedTable触发物化表刷新整体架构如下图所示从图中可以看到REST Endpoint、HiveServer2 Endpoint 等不同类型的端点共享同一个SqlGatewayService各端点负责协议接入服务层负责统一处理会话、操作与查询逻辑。快速开始环境准备SQL Gateway 被打包在 Flink 的常规发行版中因此开箱即用。它只要求存在一个可运行 Table 程序的 Flink 集群。关于集群搭建的完整说明可以参考 Cluster Deployment 文档。如果只是想快速试用可以直接用一条命令启动一个带单 worker 的本地集群$ ./bin/start-cluster.sh启动 SQL GatewaySQL Gateway 的启动脚本同样位于 Flink 的二进制发布目录bin/下源码位于 sql-gateway.sh。启动命令如下$ ./bin/sql-gateway.sh start -Dsql-gateway.endpoint.rest.addresslocalhost该命令会以 REST Endpoint 启动 SQL Gateway监听地址为localhost:8083。随后可用curl验证 REST Endpoint 是否就绪$ curl http://localhost:8083/v1/info {productName:Apache Flink,version:版本号}这里返回的productName与version信息由服务端的getGatewayInfo()能力提供。从脚本源码可以看到sql-gateway.sh完整支持四个子命令与一个帮助参数$ ./bin/sql-gateway.sh --help Usage: sql-gateway.sh [start|start-foreground|stop|stop-all] [args] commands: start - Run a SQL Gateway as a daemon start-foreground - Run a SQL Gateway as a console application stop - Stop the SQL Gateway daemon stop-all - Stop all the SQL Gateway daemons -h | --help - Show this help message其中start以守护进程方式在后台运行start-foreground以前台控制台应用方式运行便于调试和查看日志输出stop停止单个守护进程stop-all停止所有 SQL Gateway 守护进程。脚本内部对start/start-foreground会拼接FLINK_ENV_JAVA_OPTS_SQL_GATEWAY环境变量并分别委托给flink-console.sh与flink-daemon.sh执行。对于start或start-foreground命令还支持在 CLI 层面查看可配置参数$ ./bin/sql-gateway.sh start --help Start the Flink SQL Gateway as a daemon to submit Flink SQL. Syntax: start [OPTIONS] -D propertyvalue Use value for given property -h,--help Show the help message with descriptions of all options.运行 SQL 查询为了验证安装和集群连接是否正常可以按照以下三步走通第一个查询流程。整个流程基于sessionHandle会话句柄与operationHandle操作句柄两个核心标识。Step 1打开一个 Session$ curl --request POST http://localhost:8083/v1/sessions {sessionHandle:...}返回结果中的sessionHandle用于唯一标识每一个活跃用户会话。服务端会在openSession时根据客户端提供的SessionEnvironment初始化会话上下文。Step 2提交并执行一条 SQL$ curl --request POST http://localhost:8083/v1/sessions/${sessionHandle}/statements/ --data {statement: SELECT 1} {operationHandle:...}返回结果中的operationHandle用于唯一标识已提交的 SQLOperation。从 SqlGatewayService.java 的executeStatement方法签名可以看到除了语句本身还可以传入executionTimeoutMs执行超时非正值表示不启用超时机制以及executionConfig语句级执行配置。Step 3拉取结果携带上面的sessionHandle与operationHandle即可获取对应结果$ curl --request GET http://localhost:8083/v1/sessions/${sessionHandle}/operations/${operationHandle}/result/0 { results: { columns: [ { name: EXPR$0, logicalType: { type: INTEGER, nullable: false } } ], data: [ { kind: INSERT, fields: [ 1 ] } ] }, resultType: PAYLOAD, nextResultUri: ... }返回结构说明如下columns结果集的列定义包括列名与逻辑类型logicalType例如上例中EXPR$0是一个不可为空的INTEGER类型data结果数据行kind表示行类型如INSERTfields为各列的值resultType结果的传输类型例如PAYLOADnextResultUri下一页结果的拉取地址。当它不为null时说明还有后续批次的结果可以直接请求该地址获取下一批数据$ curl --request GET ${nextResultUri}通过nextResultUri逐页拉取的方式SQL Gateway 实现了大结果集的分批消费客户端无需一次性加载全部数据。SQL Gateway 配置详解SQL Gateway 启动选项如前面./bin/sql-gateway.sh --help所示脚本支持start、start-foreground、stop、stop-all与-h|--help五类命令。start命令的完整语法为start [OPTIONS]其中-D propertyvalue为给定属性设置值可多次指定-h, --help显示所有选项的帮助信息。核心配置项SQL Gateway 可以在启动时通过-Dkeyvalue方式配置也可以配置任何合法的 Flink 配置项$ ./sql-gateway -Dkeyvalue官方文档列出的 SQL Gateway 专属配置项如下表所示Key默认值类型说明sql-gateway.session.check-interval1 minDuration空闲会话超时的检查间隔设置为 0 可禁用检查。sql-gateway.session.idle-timeout10 minDuration会话在指定间隔内未被访问时的关闭超时设置为 0 表示不关闭会话。sql-gateway.session.max-num1000000IntegerSQL Gateway 服务的最大活跃会话数。sql-gateway.session.plan-cache.enabledfalseBoolean为 true 时SQL Gateway 将按会话缓存并复用查询计划。sql-gateway.session.plan-cache.size100Integer计划缓存大小仅在table.optimizer.plan-cache.enabled为 true 时生效。sql-gateway.session.plan-cache.ttl1 hourDuration计划缓存的 TTL控制缓存写入后的过期时间仅在table.optimizer.plan-cache.enabled为 true 时生效。sql-gateway.worker.keepalive-time5 minDuration空闲 worker 线程的存活时间。当 worker 数超过最小线程数时超出的线程在此间隔后被回收。sql-gateway.worker.threads.max500IntegerSQL Gateway 服务的最大 worker 线程数。sql-gateway.worker.threads.min5IntegerSQL Gateway 服务的最小 worker 线程数。这些配置项在服务实现层有着对应的落地逻辑会话生命周期管理sql-gateway.session.check-interval、sql-gateway.session.idle-timeout、sql-gateway.session.max-num由会话管理器 SessionManagerImpl.java 使用——它周期性扫描空闲会话并回收同时限制活跃会话总数对应的行为在测试 SessionManagerImplTest.java 中有覆盖计划缓存sql-gateway.session.plan-cache.*三个参数在 SessionContext.java 中读取用于控制每个会话内查询计划的缓存复用以降低重复 SQL 的解析与优化开销线程池调优sql-gateway.worker.threads.min/max与sql-gateway.worker.keepalive-time控制服务端处理请求的 worker 线程池规模与回收策略在并发请求量较大的场景下可以结合机器规格适当调整最大线程数。此外REST Endpoint 自身的监听地址与端口由 SqlGatewayRestOptions.java 定义包括Key默认值说明address无供客户端连接 SQL Gateway 服务的地址。bind-address无SQL Gateway 服务实际绑定的地址。bind-port8083服务绑定的端口支持列表如50100,50101、范围如50100-50200或两者组合多实例同机部署时建议配置端口范围避免冲突。port8083客户端连接的端口当未显式指定bind-port时服务会绑定到该端口。支持的 EndpointFlink 原生支持两类 EndpointREST EndpointSQL Gateway 默认内置的端点用户通过 HTTP/REST 接口与网关交互本文上述查询流程即基于它HiveServer2 Endpoint兼容 HiveServer2 协议的端点允许使用 Hive JDBC、Beeline 等 Hive 生态工具连接。得益于可插拔架构用户可以通过-D参数指定要启动的 Endpoint 类型$ ./bin/sql-gateway.sh start -Dsql-gateway.endpoint.typehiveserver2也可以把配置写入 Flink 配置文件sql-gateway.endpoint.type: hiveserver2注意如果 Flink 配置文件中同时存在sql-gateway.endpoint.type选项CLI 命令行参数的优先级更高。不同的 Endpoint 对应不同的接入工具REST Endpoint 面向 curl、HTTP 客户端或自研网关客户端HiveServer2 Endpoint 面向 Hive JDBC / Beeline。使用者可以根据客户端生态按需选择这正是Endpoint 决定接入方式SqlGatewayService统一处理逻辑这一架构的直接体现。总结本文围绕 SQL Gateway 的官方 Overview 文档梳理了其可插拔 Endpoint 复用型SqlGatewayService的核心架构并通过 SqlGateway.java 与 SqlGatewayService.java 印证了其装配与处理链路随后以sql-gateway.sh脚本为主线演示了启动、健康检查以及开会话 → 执行语句 → 拉取结果三步 REST 查询流程最后给出了全部会话、缓存、线程池与 REST 监听相关配置项并说明了 REST 与 HiveServer2 两种内置 Endpoint 的切换方式。读者可以依照本文流程结合 REST Endpoint 与 HiveServer2 Endpoint 两篇文档进一步深入各自端点的完整 API 细节。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink SQL Gateway 概览架构原理、REST API 实操与完整配置指南Flink SQL Gateway 概览架构原理、REST API 实操与完整配置指南 SQL Gateway 是 Apache Flink 提供的一个多客户大数据流处理批处理数据工程Flink SQL Gateway REST 端点完全指南架构、配置与 OpenAPI 实践Flink SQL Gateway REST 端点完全指南架构、配置与 OpenAPI 实践 导读 SQL Gateway 是 Flink Table 生态中大数据流处理批处理数据工程Presto Hudi Connector 完全指南架构原理、配置部署与 SQL 查询实践Presto Hudi Connector 完全指南架构原理、配置部署与 SQL 查询实践 本文面向大数据工程师与 Presto 使用者系统讲解 Prest大数据数据库后端上一篇【亲测免费】 Sony-PMCA-RE 项目推荐下一篇如何专业高效地管理Garrys Mod模组gmpublisher完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表