ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Pulsar SQL 全解析:Presto Pulsar Connector 架构、查询原理与配置实践

Apache Pulsar SQL 全解析:Presto Pulsar Connector 架构、查询原理与配置实践 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar SQL 是 Pulsar 生态中面向「结构化事件流」的 SQL 查询能力借助 Schema Registry 将 Topic 中的数据描述为带预定义字段的结构化数据再通过 Presto现称 TrinoPulsar Connector让 Presto/Trino 集群中的 Worker 直接以标准 SQL 对 Pulsar 存储的历史事件进行批式分析查询。本指南以 Pulsar 2.3.2 官方文档为骨架结合当前仓库中的连接器源码、测试用例与真实配置文件系统讲解其架构定位、基于 BookKeeper 两级分段存储的高并发读取原理、Split 拆分与谓词下推实现以及连接器全部配置项的取值与调优方向。读完本文你将能够理解 Pulsar SQL 的查询链路并掌握部署、配置与排查连接器行为所需的完整知识。Pulsar SQL 是什么在结构化事件流上直接跑 SQLApache Pulsar 本身是分布式发布订阅消息系统用于存储持续到达的事件数据流。这些事件数据并非无结构的字节流——通过 Schema Registry每个 Topic 都可以绑定一个明确的 Schema事件被写入时即按预定义字段组织为结构化数据。Pulsar SQL 正是建立在这条链路之上既然数据是结构化的就可以用关系型查询引擎对其进行 SQL 分析。在 Pulsar 中该能力由 Trino其前身为 Presto SQL承担——Presto/Trino 是一个分布式 SQL 查询引擎擅长跨数据源进行大规模交互式查询。Pulsar SQL 的整体定位如下数据面Pulsar 负责事件流的存储与读写Topic 上的 Schema 决定了可查询的字段与类型查询面Trino/Presto 集群的 Coordinator 负责任务调度与结果汇总Worker 负责实际的数据读取与计算连接面Presto Pulsar Connector 是连接两者的核心——它实现了 Trino 的 Connector SPI将 Pulsar 的 Topic 映射为可查询的表并驱动 Worker 完成数据读取。在 Pulsar 2.3.2 时代官方文档即以本仓库 sql-overview.md 作为 Pulsar SQL 的入口说明其核心论断是作为 Pulsar SQL 的核心Presto Pulsar Connector 使得 Presto 集群内的 Worker 能够查询 Pulsar 中的数据。上图展示了整体连接关系Presto 的 Coordinator 与多个 Worker 接入 Apache Pulsar 集群Worker 通过 Pulsar Connector 直接对接 Pulsar 存储层完成数据访问而不需要经过 Broker 的消费接口逐条拉取。性能基石两级分段架构与 BookKeeper 并发读取Pulsar SQL 的查询性能之所以高效且高度可扩展根源于 Pulsar 自身的 两级分段存储架构Ledger 层每个 Topic 的持久化数据在逻辑上划分为多个 LedgerSegment 层每个 Ledger 内部又由若干 SegmentEntry 序列构成Segment 是写入与读取的基本单元。Topic 中的数据以 Segment 形式存放在 Apache BookKeeper 中而不是只存在单一副本。每个 Topic Segment 都会被复制到若干 BookKeeper 节点上这种多副本布局直接带来了两个关键收益并发读取同一段数据可以在多个副本节点上并行读取读吞吐不受单点瓶颈限制水平扩展副本节点数可通过配置调整BookKeeper 节点数的默认值为3副本越多可供并发读取的节点就越多。对 Pulsar SQL 而言最核心的一点在于Presto Pulsar Connector 直接从 BookKeeper 读取数据。这意味着查询流量不经过 Pulsar Broker 的消费链路而是由 Presto Worker 以只读游标直连 BookKeeper 存储层从而可以在水平可扩展的 BookKeeper 节点集合上实现并发读取。这是 Pulsar SQL 与「从 Pulsar 客户端逐条消费再做分析」的本质区别。上图展示了查询执行的具体形态Presto Coordinator 统筹调度多个 Presto Worker 将 Pulsar 事件流拆分为逻辑 Segment 并各自加载对应区间并行读取事件流中的 Latest Reader 与 Late Reader 分别映射到不同时间位置的 Segment 数据体现出「读最新数据」与「回溯读历史数据」两种查询都能被同一套分段机制高效承载。源码视角Connector 如何把 Topic 拆成 Split 并行读取连接器的工作方式可以用「一个查询 多个 Split一个 Split 一段 Entry 区间」来概括。相关实现集中在 pulsar-sql/presto-pulsar 模块。Split 的生成PulsarSplitManagerPulsarSplitManager.java 实现了 Trino 的ConnectorSplitManager接口其getSplits()方法完成查询计划到读取任务的关键转换通过PulsarAdmin从 Pulsar 获取目标 Topic 的 Schema 信息pulsarAdmin.schemas().getSchemaInfo(...)若未定义 Schema 则回退到默认 Schema判断 Topic 是否为分区 Topic非分区 Topic 走getSplitsNonPartitionedTopic分区 Topic 走getSplitsPartitionedTopic无论哪种路径最终都会进入getSplitsForTopic其核心操作是用managedLedgerFactory.openReadOnlyCursor(...)在 Topic 对应的 Managed Ledger 上打开一个只读游标读取总 Entry 数numEntries按配置的target-num-splits计算每个 Split 应承担的 Entry 数余数均摊得到形如(startPosition, endPosition)的 Entry 区间列表每个区间封装为一个PulsarSplit返回给 Trino 调度器由不同 Worker 并行处理。谓词下推只读需要的数据getSplitsForTopic中还有一个值得注意的优化——PredicatePushdownInfo.getPredicatePushdownInfo(...)当查询的 WHERE 条件命中内部列__publish_time__且为单个时间范围时连接器会在拆分阶段就基于时间范围定位起始与结束 Entry 位置findPosition通过ReadOnlyCursor.findNewestMatching按 Entry 时间戳二分定位从而将 Split 裁剪到最小所需区间当查询命中内部列__partition__时getPredicatedPartitions(...)会只对满足分区条件的分区生成 Split减少不必要的扫描。这种「先裁剪、再读取」的下推机制使 Pulsar SQL 在处理时间范围或分区过滤查询时可以避免扫描整个 Topic 的全部数据。记录读取PulsarRecordCursor 的流水线每个 Split 由一个 PulsarRecordCursor.java 消费它实现了 Trino 的RecordCursor接口。从源码结构看其读取是一个两级队列流水线通过ReadOnlyCursor从 BookKeeper 读取原始 Entry进入SpscArrayQueueEntryentryQueue后台线程对 Entry 做反序列化得到RawMessage放入SpscArrayQueueRawMessagemessageQueue消费侧从 messageQueue 弹出消息按列句柄PulsarColumnHandle逐字段解码为 Trino 的行值同时追踪completedBytes、entriesProcessed等统计指标。两个队列的大小分别受配置项pulsar.max-split-entry-queue-size与pulsar.max-split-message-queue-size约束整体缓存上限由pulsar.max-split-queue-cache-size控制相关的内存分配逻辑位于 util/CacheSizeAllocator.java 及其两个变体NoStrictCacheSizeAllocator、NullCacheSizeAllocator中。整个模块的拆分与读取行为有对应的单元测试覆盖可参考 TestPulsarSplitManager.java 与 TestPulsarRecordCursor.java。字段模型Schema 如何映射为可查询的列Connector 将 Topic 映射为表后列的来源有两类Schema 中定义的业务字段以及连接器内置的元数据列。多格式行解码器decoder 目录按消息编码格式提供了四类行解码器由PulsarDispatchingRowDecoderFactory依据 Topic Schema 类型分发解码器适用 Schema 类型关键实现Avro 解码器AVROPulsarAvroRowDecoderFactory/PulsarAvroColumnDecoderJSON 解码器JSONPulsarJsonRowDecoderFactory/PulsarJsonFieldDecoderProtobufNative 解码器PROTOBUF_NATIVEPulsarProtobufNativeRowDecoderFactory/PulsarProtobufNativeColumnDecoderPrimitive 解码器基础类型STRING 等PulsarPrimitiveRowDecoderFactory/PulsarPrimitiveRowDecoder每种解码器均有对应测试类验证其字段解析正确性如 TestAvroDecoder.java、TestJsonDecoder.java、TestProtobufNativeDecoder.java 与 TestPrimitiveDecoder.java。内置元数据列除业务字段外连接器还为每行附加了内部列定义见 PulsarInternalColumn.java。这些列可用于 WHERE 过滤与排序也是谓词下推的入口内部列类型含义__partition__INTEGER消息所属分区号__event_time__TIMESTAMP应用定义的事件发生时间戳毫秒__publish_time__TIMESTAMP消息发布的时间戳毫秒__message_id__VARCHAR生成该行的消息 ID__sequence_id__BIGINT生成该行的消息序列号__producer_name__VARCHAR发布该消息的生产者名称__key__VARCHARTopic 的分区键__properties__VARCHAR用户自定义属性连接器配置详解从默认值到调优连接器的所有配置项在 PulsarConnectorConfig.java 中声明每个配置项对应Config(pulsar.xxx)注解仓库同时提供了可直接使用的真实配置文件 conf/presto/catalog/pulsar.properties将其放入 Trino/Presto 的catalog目录即可加载 Pulsar 目录。以下按类别给出完整参数说明。连接与 Schema 访问配置项默认值说明connector.namepulsarCatalog 中展示的连接器名称必须为pulsarpulsar.broker-service-urlhttp://localhost:8080Pulsar Broker 服务地址已标记 DEPRECATEDpulsar.web-service-urlhttp://localhost:8080Pulsar Broker Web 服务地址从源码看若配置了该值连接器优先使用它构造PulsarAdminpulsar.zookeeper-urilocalhost:2181Zookeeper 集群地址用于获取 BookKeeper 元数据定位读取与拆分性能配置项默认值说明pulsar.max-entry-read-batch-size100单次批量读取的最小 Entry 数影响读放大与批处理效率pulsar.target-num-splits2每次查询默认生成的 Split 数数值越大并行度越高但每个 Split 越小pulsar.max-split-message-queue-size10000每个 Split 的消息队列最大条数pulsar.max-split-entry-queue-size1000每个 Split 的 Entry 队列最大条数pulsar.max-split-queue-cache-size-1Split 队列缓存字节上限该值的一半用于 Entry 队列字节、一半用于消息队列字节-1表示不限制pulsar.max-message-size52428805MB单条批量消息的最大尺寸命名空间与分区配置项默认值说明pulsar.namespace-delimiter-rewrite-enablefalse是否重写命名空间分隔符pulsar.rewrite-namespace-delimiter/重写使用的分隔符注意其值不能与命名空间允许的字符a-zA-Z_0-9 -:%冲突否则连接器会抛出异常认证与 TLS配置项说明pulsar.auth-plugin用于认证 Pulsar 集群的认证插件类名pulsar.auth-params认证参数pulsar.tls-allow-insecure-connection是否接受不受信任的 TLS 证书pulsar.tls-hostname-verification-enable是否对 TLS 连接启用主机名校验pulsar.tls-trust-cert-file-path受信任 TLS 证书文件路径从 PulsarConnectorConfig.getPulsarAdmin() 的实现可以看到认证与 TLS 配置最终会透传给PulsarAdmin.builder()的authentication、allowTlsInsecureConnection、enableTlsHostnameVerification、tlsTrustCertsFilePath等方法用于构造访问 Pulsar 元数据的 Admin 客户端。BookKeeper 客户端配置项默认值说明pulsar.bookkeeper-throttle-value0每秒钟 Entry 读取的限流阈值0表示禁用限流pulsar.bookkeeper-num-io-threads2 × CPU 核数Netty 处理 TCP 连接的线程数pulsar.bookkeeper-num-worker-threadsCPU 核数BookKeeper 客户端提交操作的 worker 线程数pulsar.bookkeeper-use-v2-protocoltrue是否使用 v2 协议LAC 由读指针携带false时使用 v3 协议并启用显式 LACpulsar.bookkeeper-explicit-interval0v3 协议下显式 LAC 的发送间隔Managed Ledger 与分层存储配置项默认值说明pulsar.managed-ledger-cache-size-MB0Managed Ledger 数据缓存大小MB从 JVM 直接内存分配同一 SQL Worker 内共享0表示禁用缓存pulsar.managed-ledger-num-worker-threadsCPU 核数Managed Ledger 任务分发线程数pulsar.managed-ledger-num-scheduler-threadsCPU 核数Managed Ledger 定时任务线程数pulsar.managed-ledger-offload-driver空将旧数据卸载到长期存储的驱动如aws-s3pulsar.offloaders-directory./offloaders存放 Offloader 的目录pulsar.managed-ledger-offload-max-threads2Ledger 卸载线程池的最大线程数pulsar.offloader-properties空具体 Offloader 实现的属性JSON 格式例如 S3 的桶、区域与端点监控配置项说明pulsar.stats-provider统计提供器类名默认org.apache.bookkeeper.stats.NullStatsProvider可切换为 Prometheus 提供器pulsar.stats-provider-configs统计提供器的 JSON 配置如 Prometheus 的 HTTP 端口与启用开关配置文件示例可直接复制使用# name of the connector to be displayed in the catalog connector.namepulsar # the url of Pulsar broker service # DEPRECATED pulsar.broker-service-urlhttp://localhost:8080 # the url of Pulsar broker web service pulsar.web-service-urlhttp://localhost:8080 # URI of Zookeeper cluster pulsar.zookeeper-uri127.0.0.1:2181 # minimum number of entries to read at a single time pulsar.max-entry-read-batch-size100 # default number of splits to use per query pulsar.target-num-splits2 # max message queue size pulsar.max-split-message-queue-size10000 # max entry queue size pulsar.max-split-entry-queue-size1000 # half of this value is used as max entry queue size bytes and the left is used as max message queue size bytes pulsar.max-split-queue-cache-size-1部署集成仓库中的相关模块Pulsar SQL 在仓库中由pulsar-sql目录下的多模块构成构建后可获得独立的分发物presto-pulsar连接器核心实现包含上述全部 Connector SPI 类与解码器其 PulsarPlugin.java 是 Trino 插件入口presto-pulsar-plugin连接器的打包assembly模块负责将连接器及其依赖组装为可放入 Trino 插件目录的 NAR/JAR 包presto-distributionPresto/Trino 发行版组装模块产出可直接运行的 SQL 引擎分发物java-version-trim-agent一个 JVM 代理工具TrimJavaVersionAgent用于在运行期对 Java 版本字符串做裁剪处理部分 Presto 组件对 JDK 版本的兼容性校验问题。部署时的一般做法是将presto-pulsar-plugin构建出的插件放入 Presto/Trino 的plugin目录将上述catalog/pulsar.properties放入etc/catalog目录随后即可在 SQL 客户端中通过pulsar.tenant/namespace/topic形式的三段式命名访问 Topic 数据。小结Pulsar SQL 的核心价值在于把「消息系统里按序到达的事件流」变成「可以被标准 SQL 直接查询的分布式表」。支撑这一能力的四根支柱环环相扣Schema Registry 提供结构化字段描述Presto Pulsar Connector 完成 Topic 到表的映射与查询拆分两级分段存储与 BookKeeper 多副本提供并发读取的物理基础而 Split 裁剪与谓词下推则保证查询只读取必要的数据区间。对于需要针对 Pulsar 历史事件做分析、报表或回溯检索的场景Pulsar SQL 提供了一条无需导出数据、直接在存储层完成计算的查询路径。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar SQL 部署与 Presto Pulsar Connector 配置实战指南Apache Pulsar SQL 部署与 Presto Pulsar Connector 配置实战指南 本文以 Apache Pulsar 2.2.1 版本官消息队列后端流处理Apache Pulsar SQLPresto Pulsar Connector部署与配置完整指南Apache Pulsar SQLPresto Pulsar Connector部署与配置完整指南 Apache Pulsar SQL 是构建在 Prest消息队列后端流处理Apache Pulsar SQL 架构解析基于 Presto Connector 的分布式流式数据查询Apache Pulsar SQL 架构解析基于 Presto Connector 的分布式流式数据查询 Pulsar 最典型的应用场景之一是存储结构化事件数消息队列后端流处理上一篇洛雪音乐音源终极指南免费解锁全网无损音乐的完整解决方案下一篇5分钟掌握智慧教育平台电子课本高效下载方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表