Apache Arrow Flight SQL 协议详解RPC 命令、执行模型与会话管理【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow12/arrowFlight SQL 是 Apache Arrow 定义的一套标准协议用于通过 Arrow 内存格式与 Flight RPC 框架访问 SQL 数据库其权威规范位于 docs/source/format/FlightSql.rst完整的 Protobuf 定义见 format/FlightSql.proto。本文从协议设计出发逐条梳理 SQL 元数据命令、查询执行命令、会话管理命令的用法与底层调用链并配合仓库中的序列图、Java/Go 客户端实现帮助你理解数据库只需要实现 RPC 服务端、客户端即可复用通用驱动这一核心架构掌握实现或对接 Flight SQL 服务所需的关键细节。Flight SQL 的定位与整体设计Flight SQL 本质上是一个客户端-服务端协议约定数据库厂商按规范实现一组 RPC 方法服务端而客户端无需为每个数据库单独编写驱动。任何支持相应端点的数据库都可以被一个通用的 Flight SQL 客户端访问。Flight SQL 客户端包装底层的 Flight 客户端为下面描述的新 RPC 方法提供便捷 API。从 FlightSql.rst 可以提炼出协议的三条基本规则复用 Flight 既有 RPCFlight SQL 不发明新的传输协议而是复用 Arrow Flight 预定义的GetFlightInfo、GetSchema、DoGet、DoPut、DoAction等方法命令驱动每条命令是一个通过 Protobuf 定义的请求/响应消息与某个 Flight RPC 方法配对使用结果即 Arrow 数据SQL 元数据与查询结果都以 Arrow 数据的形式返回schema 由规范固定客户端可以据此直接消费。命令的编码方式所有命令类消息Command*前缀统一编码进FlightDescriptor的cmd字段CMD 类型描述符先把 Protobuf 请求消息打包进google.protobuf.Any序列化后填入cmd字段。这条规则在 FlightSql.proto 的几乎所有命令注释中都有明示例如CommandStatementQuery的注释写明其用于 GetSchema/GetFlightInfo且Fields on this schema may contain the following metadata——即返回的 schema 字段上可以携带ARROW:FLIGHT:SQL:*元数据详见下文。SQL 元数据命令让客户端看懂数据库Flight SQL 提供一组命令用于获取数据库服务端的目录catalog元数据。它们全部可以与GetFlightInfo和GetSchema两个 RPC 方法配对使用命令作用CommandGetCatalogs列出数据库中可用的 catalog定义因厂商而异CommandGetCrossReference列出给定表中引用父表列的外键列CommandGetDbSchemas列出数据库中可用的 schema注意是表的分组不是 Arrow schema定义因厂商而异CommandGetExportedKeys列出引用给定表主键列的外键列CommandGetImportedKeys列出给定表的外键CommandGetPrimaryKeys列出给定表的主键CommandGetSqlInfo获取数据库服务端及其 SQL 功能支持的元数据CommandGetTables列出数据库中的表CommandGetTableTypes列出数据库中的表类型类型列表因厂商而异元数据命令的调用流程规范明确规定若命令与GetFlightInfo一起使用服务端返回FlightInfo客户端随后用FlightInfo中的 Ticket 调用DoGet获取包含命令结果的 Arrow 数据——SQL 元数据和查询结果一样以 Arrow 数据返回特定命令由GetSchema或DoGet返回的 Arrow schema 是规范固定的。仓库中的序列图 docs/source/format/FlightSql/CommandGetTables.mmd 清晰展示了这一流程Client-Server: GetFlightInfo(CommandGetTables) Server-Client: FlightInfo{…Ticket…} Client-Server: DoGet(Ticket) Server-Client: stream of FlightData关键命令的返回 schema取自 FlightSql.proto以下 schema 均为规范固定客户端可直接据此解析结果CommandGetCatalogs单列catalog_name: utf8 not null按catalog_name排序CommandGetDbSchemascatalog_name: utf8db_schema_name: utf8 not null支持catalog与db_schema_filter_pattern过滤参数模式串中%匹配任意子串、_匹配单个字符CommandGetTablescatalog_name、db_schema_name、table_name、table_type均 utf8其中 table_name/table_type 非空可选table_schema: bytes按 Schema.fbs 序列化的 IPC 消息支持catalog、db_schema_filter_pattern、table_name_filter_pattern、table_types过滤及include_schema开关CommandGetTableTypes单列table_type: utf8 not null常见值如TABLE、VIEW、SYSTEM TABLECommandGetPrimaryKeyscatalog_name、db_schema_name、table_name、column_name、key_nameutf8key_sequence: int32CommandGetExportedKeys/CommandGetImportedKeys/CommandGetCrossReference返回主键表与外键表的 catalog/schema/table/column、key_sequence、update_rule/delete_ruleuint8对应UpdateDeleteRules枚举0CASCADE、1RESTRICT、2SET NULL、3NO ACTION、4SET DEFAULT。CommandGetSqlInfo 与 SqlInfo 枚举CommandGetSqlInfo是元数据体系中最丰富的一个其返回 schema 为info_name: uint32 not null, value: dense_union string_value: utf8, bool_value: bool, bigint_value: int64, int32_bitmask: int32, string_list: liststring_data: utf8, int32_to_int32_list_map: mapkey: int32, value: listint32 info字段为repeated uint32省略时返回全部元数据。FlightSql.proto 中的SqlInfo枚举对信息 ID 做了严格分区[0-500) 服务器信息FLIGHT_SQL_SERVER_NAME(0)、FLIGHT_SQL_SERVER_VERSION(1)、FLIGHT_SQL_SERVER_ARROW_VERSION(2)、FLIGHT_SQL_SERVER_READ_ONLY(3)、FLIGHT_SQL_SERVER_SQL(4)、FLIGHT_SQL_SERVER_SUBSTRAIT(5) 及 Substrait 版本范围(6/7)、FLIGHT_SQL_SERVER_TRANSACTION(8取值见SqlSupportedTransaction)、FLIGHT_SQL_SERVER_CANCEL(9)、FLIGHT_SQL_SERVER_BULK_INGESTION(10)、FLIGHT_SQL_SERVER_INGEST_TRANSACTIONS_SUPPORTED(11)、FLIGHT_SQL_SERVER_STATEMENT_TIMEOUT(100)、FLIGHT_SQL_SERVER_TRANSACTION_TIMEOUT(101)[500-1000) SQL 语法信息从SQL_DDL_CATALOG(500) 一直到SQL_STORED_FUNCTIONS_USING_CALL_SYNTAX_SUPPORTED(576)涵盖大小写敏感性、标识符引用符、关键字/数值函数/字符串函数/系统函数/日期时间函数列表、通配符转义、GROUP BY 支持位掩码SQL_SUPPORTED_GROUP_BY、SQL 语法等级位掩码SQL_SUPPORTED_GRAMMAR、ANSI92 等级、外连接支持等级、事务隔离级别与支持位掩码、结果集类型与并发性位掩码、类型转换映射SQL_SUPPORTS_CONVERT返回mapint32, listint32等。SqlInfo的取值大量借鉴 ODBC 的SQLGetInfo()函数设计[0, 10000)区间预留给规范默认项自定义元数据 ID 应从 10000 开始。服务端至少应返回规范指定集合也可以返回更多。CommandGetXdbcTypeInfoCommandGetXdbcTypeInfo可选data_type过滤参数返回数据类型的详细描述返回列包括type_name、data_type值同 JDBC/ODBC 的XdbcDataType枚举、column_size、literal_prefix/literal_suffix、create_params、nullable见Nullable枚举NO_NULLS/NULLABLE/UNKNOWN、case_sensitive、searchable见Searchable枚举、unsigned_attribute、fixed_prec_scale、auto_increment、local_type_name、minimum_scale/maximum_scale、sql_data_type、datetime_subcode见XdbcDatetimeSubcode、num_prec_radix、interval_precision等。查询执行命令即席查询、更新与批量摄取除元数据外Flight SQL 还提供执行 SQL 查询与管理预编译语句prepared statement的命令。其中多数命令与元数据命令一样通过GetFlightInfo/GetSchema使用部分命令可与DoPut配对此时命令仍编码在请求FlightDescriptor中。即席查询CommandStatementQuery / CommandStatementUpdateCommandStatementQuery执行一次即席 SQL 查询。与GetFlightInfo配合执行查询随后DoGet拉取结果与GetSchema配合返回查询结果的 schema。消息字段query: string必填 可选的transaction_id: bytes未设置则自动提交。CommandStatementUpdate执行不返回结果的即席 SQL如 DML/DDL。与DoPut配合执行查询并返回受影响行数。CommandStatementQuery的完整调用链见序列图 docs/source/format/FlightSql/CommandStatementQuery.mmd注意GetFlightInfo返回的FlightInfo可能包含多个FlightEndpoint客户端需要遍历每个 endpoint依次DoGet(endpoint.ticket)每个 ticket 返回一段 FlightData 流。更新类命令的响应DoPutUpdateResult执行更新的命令CommandStatementUpdate、CommandStatementIngest在消费完整个 FlightData 流后返回一个 Flight SQL 的DoPutUpdateResult。该消息被编码在 Flight RPCPutResult的app_metadata字段中返回。DoPutUpdateResult的定义为message DoPutUpdateResult { // The number of records updated. A return value of -1 represents // an unknown updated record count. int64 record_count 1; }即record_count为受影响行数-1表示行数未知。预编译语句创建、绑定、执行、关闭预编译语句的生命周期涉及三类 RPC创建ActionCreatePreparedStatementRequest携带query: string与可选transaction_id。以Action开头的命令通过DoActionRPC 使用命令打包进google.protobuf.Any后序列化进 Flight Action 的body同时Action 的type字段必须设为命令名——例如ActionCreatePreparedStatementRequest的type应为CreatePreparedStatement。响应ActionCreatePreparedStatementResult包含prepared_statement_handle不透明句柄用于标识预编译语句、可选的dataset_schema结果集 schemaIPC 封装与parameter_schema绑定参数 schemaIPC 封装注意结果集 schema 可能依赖绑定参数因此这里返回的 schema 未必准确客户端不应假设它等于实际执行返回数据的 schema对于无具体类型的绑定参数如SELECT ?规范未规定如何处理建议用 union 类型枚举可能类型或用 NAnull类型作为通配符/占位。绑定与执行CommandPreparedStatementQuery与DoPut配合将参数值绑定到预编译语句与GetFlightInfo配合执行预编译语句语句可在取回结果后复用与GetSchema配合获取结果集的预期 schema若之前已用 DoPut 绑定过参数服务端应把这些值纳入考量。句柄更新机制DoPut之后服务端可能返回更新后的句柄DoPutPreparedStatementResult.prepared_statement_handle可选。更新句柄允许无状态服务把已绑定的参数编码进新句柄客户端后续请求必须使用新句柄上一次 DoPut 返回的句柄还可以再次传给下一次CommandPreparedStatementQuery的 DoPut 以绑定新参数集服务端负责检测客户端未使用更新句柄的情况并返回错误。更新执行CommandPreparedStatementUpdate—— 执行不返回结果的预编译语句与DoPut配合返回受影响行数语句之后可复用。关闭ActionClosePreparedStatementRequesttypeClosePreparedStatement携带prepared_statement_handle关闭服务端与该句柄关联的资源。句柄也可由服务端超时自动回收见FLIGHT_SQL_SERVER_STATEMENT_TIMEOUT。完整的生命周期序列见 docs/source/format/FlightSql/CommandPreparedStatementQuery.mmdDoAction(CreatePreparedStatement)→ 循环DoPut绑定参数 → 可选更新句柄 →GetFlightInfo执行 → 遍历 endpointsDoGet→DoAction(ClosePreparedStatement)。批量摄取CommandStatementIngestCommandStatementIngest用于批量加载数据与DoPut配合把 Arrow record batch 流加载进指定目标表并通过DoPutUpdateResult返回摄取行数。序列见 docs/source/format/FlightSql/CommandStatementIngest.mmd。其消息字段包括message CommandStatementIngest { message TableDefinitionOptions { enum TableNotExistOption { TABLE_NOT_EXIST_OPTION_UNSPECIFIED 0; // 客户端不得使用 TABLE_NOT_EXIST_OPTION_CREATE 1; // 表不存在则创建 TABLE_NOT_EXIST_OPTION_FAIL 2; // 表不存在则失败 } enum TableExistsOption { TABLE_EXISTS_OPTION_UNSPECIFIED 0; // 客户端不得使用 TABLE_EXISTS_OPTION_FAIL 1; // 表已存在则失败 TABLE_EXISTS_OPTION_APPEND 2; // 表已存在则追加 TABLE_EXISTS_OPTION_REPLACE 3; // 表已存在则重建 } TableNotExistOption if_not_exist 1; TableExistsOption if_exists 2; } TableDefinitionOptions table_definition_options 1; string table 2; optional string schema 3; optional string catalog 4; bool temporary 5; // 临时表后端定义命名空间会话结束即删除 optional bytes transaction_id 6; mapstring, string options 1000; // 后端特定选项字段号 1000 以下预留给规范扩展 }Substrait 与事务相关动作协议还包含若干扩展能力Substrait 支持ActionCreatePreparedSubstraitPlanRequest与CommandStatementSubstraitPlan允许执行序列化的 Substrait 计划SubstraitPlan消息携带plan: bytes与version: string如0.12.0用于告知消费方计划版本。事务动作BeginTransaction、EndTransactionCOMMIT/ROLLBACK、BeginSavepoint、EndSavepointRELEASE/ROLLBACK TO句柄均为不透明 bytes保存点仅当服务端事务支持级别为 SAVEPOINT 时可用。查询取消ActionCancelQueryRequest携带序列化 FlightInfo与ActionCancelQueryResult用于显式取消运行中的查询。注意自 13.0.0 起该命令已废弃改用 Flight 的CancelFlightInfo动作option deprecated true标注在 FlightSql.proto 中可见。Flight Server 会话管理SetSessionOptions / GetSessionOptions / CloseSessionFlight SQL 提供会话级命令用于设置和更新影响服务端行为的会话变量。常见选项取决于服务端实现包括catalog与schema用于指示查询作用的当前 catalog 与 schema。SetSessionOptions按名称/值设置服务端会话选项GetSessionOptions获取当前会话选项包括客户端设置的以及服务端默认/隐式设置的CloseSession关闭并使当前会话上下文失效。规范给出的实践建议客户端应尽量在发起查询和其他命令之前设置选项因为某些服务端实现要求选项恰好设置一次且必须先于任何可能触发隐式设置的其它活动为兼容 JDBC/ODBC 等 Database Connectivity 驱动强烈建议服务端接受所有选项值的字符串表示驱动可能把连接串中的选项原样透传给服务端而不做转换也建议接受并转换其它数值类型为目标类型但非强制会话通过实现定义的方式在客户端与服务端之间持久化通常是 RFC 6265 cookie服务端也可把其它连接状态不透明地并入会话令牌。CloseSession会使通过会话上下文持久化的任何认证上下文一并失效会话可在非空或空的SetSessionOptions调用时启动或由服务端自行选择时机。客户端/服务端实现参考仓库中的落地代码协议在仓库的多个语言绑定中均有实现可作为阅读与对接的参照Javajava/flight/flight-sql 模块提供官方实现核心类 FlightSqlClient.java 暴露了与本文命令一一对应的 APIgetCatalogs()、getSchemas()、getTables()、getSqlInfo()、getPrimaryKeys()、getExportedKeys()、getImportedKeys()、getCrossReference()、getTableTypes()、getXdbcTypeInfo()、execute()、executeUpdate()、prepare()以及各自的*Schema()变体服务端侧可继承 FlightSqlProducer.java 或空实现 NoOpFlightSqlProducer.java 快速起步示例服务端见 FlightSqlExample.java集成测试场景见 FlightSqlScenarioProducer.javaGogo/arrow/flight/flightsql 提供Clientclient.go与Serverserver.go并附带完整的 SQLite 示例服务端 sqlite_server.go 与测试 client_test.goC GLibc_glib/arrow-flight-sql-glib 提供 C 语言绑定ruby/red-arrow-flight-sql 提供 Ruby 绑定其测试位于 c_glib/test/flight-sql协议生成代码Java 端由arrow.flight.protocol.sql包提供Go 端生成代码在 FlightSql.pb.go。总结Flight SQL 通过复用 Flight RPC Protobuf 命令 Arrow 数据作为统一结果载体三层设计把数据库访问标准化服务端只需按 FlightSql.proto 实现命令端点客户端即可用统一的 Flight SQL 驱动完成元数据浏览、即席查询、预编译语句、批量摄取与会话管理。元数据命令统一走GetFlightInfo/GetSchemaDoGet取 Arrow 流Action*命令走DoActiontype为命令名更新类命令经DoPut后以DoPutUpdateResult返回行数——抓住这几条主线即可快速读懂并实现一个符合规范的 Flight SQL 服务。【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow12/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
