- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文是 Apache Pulsar 中Pulsar SQL(由 Presto/Trino 引擎驱动的内置 SQL 查询能力)的完整上手指南。你将学习如何在本地 standalone 环境中依次启动 Pulsar 集群与 SQL worker、进入pulsar sqlCLI,并通过内置data-generator连接器注入模拟数据后用SELECT查询;最后掌握用 Java Producer + Avro Schema 写入自定义数据并在 SQL 中查询的完整链路。读完本文,你可以独立搭建一套可用的 Pulsar SQL 开发环境,并理解其底层读取架构。本文内容基于仓库site2/website-next/versioned_docs/version-2.3.2/sql-getting-started.md编写,并补充了对应源码与配置佐证。
前置条件与准备工作
在开始查询 Pulsar 中的数据之前,需要先完成两件事(对应 Requirements 一节):
- 安装 Pulsar standalone:参考 Set up a standalone Pulsar locally。standalone 模式会将 Pulsar broker、必要的 ZooKeeper 与 BookKeeper 组件运行在同一个 JVM 进程中,是本地开发与测试的最简形态。
- 安装 Pulsar 内置连接器(built-in connectors):参考 Install builtin connectors (optional)。Pulsar SQL 快速入门需要用到
data-generator这个内置 Source 连接器来注入测试数据,因此这一步不能省略。自2.1.0-incubating起,内置连接器以单独的二进制分发包发布,需要将其中的.nar文件拷贝到 Pulsar 目录下的connectors目录中。
说明:本指南以仓库中 version-2.3.2 的文档为准,涉及的 CLI 命令(
bin/pulsar standalone、bin/pulsar sql-worker run、bin/pulsar sql)在 2.3.x 系列版本中一致。
三步启动:standalone 集群 + SQL worker + SQL CLI
在满足上述条件后,按以下顺序启动三个进程。每一条命令都应在 Pulsar 解压目录下执行。
第 1 步:启动 Pulsar standalone 集群。
./bin/pulsar standalone启动成功后,会在当前终端持续输出日志。由于该服务占用当前终端,后续命令请另开新的终端窗口执行。
第 2 步:启动 Pulsar SQL worker。
./bin/pulsar sql-worker runPulsar SQL 的查询能力由 Trino(前身为 Presto SQL) 提供,sql-worker命令本质上是 Presto launcher 的封装,支持run、start、stop、restart、kill、status等子命令。run表示在前台运行;如需后台守护进程方式,可使用./bin/pulsar sql-worker start。
第 3 步:启动 SQL CLI。
./bin/pulsar sql等待 standalone 集群与 SQL worker 初始化完成后,CLI 会进入presto>交互式命令行提示符,此时即可输入 SQL 语句。
第一条 SQL:验证 Pulsar SQL 环境就绪
进入presto>提示符后,依次执行以下命令,验证集群状态与目录结构。
查看 Catalog(目录):
presto> show catalogs; Catalog --------- pulsar system (2 rows) Query 20180829_211752_00004_7qpwh, FINISHED, 1 node Splits: 19 total, 19 done (100.00%) 0:00 [0 rows, 0B] [0 rows/s, 0B/s]pulsarcatalog 由 Presto Pulsar connector 注册(配置项为connector.name=pulsar),system是 Presto 引擎自带的内置 catalog,用于查询集群节点信息(例如SELECT * FROM system.runtime.nodes)。
查看 pulsar catalog 下的 Schema(命名空间):
presto> show schemas in pulsar; Schema ----------------------- information_schema public/default public/functions sample/standalone/ns1 (4 rows) Query 20180829_211818_00005_7qpwh, FINISHED, 1 node Splits: 19 total, 19 done (100.00%) 0:00 [4 rows, 89B] [21 rows/s, 471B/s]Schema 与 Pulsar 的命名空间一一对应:public/default是 standalone 启动时自动创建的默认开发命名空间,sample/standalone/ns1同样由 standalone 模式预置,information_schema则是由 Presto 提供的元数据 Schema。
查看某个 Schema 下有哪些表(Topic):
presto> show tables in pulsar."public/default"; Table ------- (0 rows) Query 20180829_211839_00006_7qpwh, FINISHED, 1 node Splits: 19 total, 19 done (100.00%) 0:00 [0 rows, 0B] [0 rows/s, 0B/s]此时 Pulsar 中还没有任何数据,所以返回 0 行。在 Pulsar SQL 的映射模型中,一个带 Schema 的 Topic 就等价于一张表,目录/表名的三层结构为catalog."tenant/namespace".table_name。
用内置>./bin/pulsar-admin sources create --name generator --destinationTopicName generator_test --source-type>presto> show tables in pulsar."public/default"; Table ---------------- generator_test (1 row) Query 20180829_213202_00000_csyeu, FINISHED, 1 node Splits: 19 total, 19 done (100.00%) 0:02 [1 rows, 38B] [0 rows/s, 17B/s]
第 6 步:查询 Topic 中的数据。
presto> select * from pulsar."public/default".generator_test; firstname | middlename | lastname | email | username | password | telephonenumber | age | companyemail | nationalidentitycardnumber | -------------+-------------+-------------+----------------------------------+--------------+----------+-----------------+-----+-----------------------------------------------+----------------------------+ Genesis | Katherine | Wiley | genesis.wiley@gmail.com | genesisw | y9D2dtU3 | 959-197-1860 | 71 | genesis.wiley@interdemconsulting.eu | 880-58-9247 | Brayden | | Stanton | brayden.stanton@yahoo.com | braydens | ZnjmhXik | 220-027-867 | 81 | brayden.stanton@supermemo.eu | 604-60-7069 | Benjamin | Julian | Velasquez | benjamin.velasquez@yahoo.com | benjaminv | 8Bc7m3eb | 298-377-0062 | 21 | benjamin.velasquez@hostesltd.biz | 213-32-5882 | Michael | Thomas | Donovan | donovan@mail.com | michaeld | OqBm9MLs | 078-134-4685 | 55 | michael.donovan@memortech.eu | 443-30-3442 | Brooklyn | Avery | Roach | brooklynroach@yahoo.com | broach | IxtBLafO | 387-786-2998 | 68 | brooklyn.roach@warst.biz | 085-88-3973 | Skylar | | Bradshaw | skylarbradshaw@yahoo.com | skylarb | p6eC6cKy | 210-872-608 | 96 | skylar.bradshaw@flyhigh.eu | 453-46-0334 | . . .DataGeneratorSource 的实现机制
从仓库源码可以印证上述命令背后发生了什么。DataGeneratorSource位于 pulsar-io/data-generator/src/main/java/org/apache/pulsar/io/datagenerator/DataGeneratorSource.java,其核心逻辑是:
- 在
open()中加载配置并创建 Fairy 实例——一个用于生成随机个人信息的 Java 库; - 在
read()中每次Thread.sleep(...)后返回一条Person记录(fairy.person()),字段包括 firstname、lastname、email、username、password、age 等,与上面查询结果的列完全对应; - 消息之间的发送间隔由 DataGeneratorSourceConfig.java 中的
sleepBetweenMessages控制,默认值为50(毫秒),标注为@PositiveNumber校验。
因此该连接器会以大约每 50ms 一条的速度,源源不断地把结构化的Person数据写入generator_testTopic,Pulsar SQL 将其识别为一张结构化表,支持任意SELECT查询。
查询你自己的数据:Java Producer + Avro Schema
如果不想使用模拟数据,可以先把自己的数据写入 Pulsar,再通过 Pulsar SQL 查询。要点是消息必须携带 Schema——Pulsar SQL 依赖 Schema Registry 来将 Topic 映射为带列结构的表。以下是一个使用 Avro Schema 的 Java Producer 完整示例(原文档示例,可直接复制运行):
public class TestProducer { public static class Foo { private int field1 = 1; private String field2; private long field3; public Foo() { } public int getField1() { return field1; } public void setField1(int field1) { this.field1 = field1; } public String getField2() { return field2; } public void setField2(String field2) { this.field2 = field2; } public long getField3() { return field3; } public void setField3(long field3) { this.field3 = field3; } } public static void main(String[] args) throws Exception { PulsarClient pulsarClient = PulsarClient.builder().serviceUrl("pulsar://localhost:6650").build(); Producer<Foo> producer = pulsarClient.newProducer(AvroSchema.of(Foo.class)).topic("test_topic").create(); for (int i = 0; i < 1000; i++) { Foo foo = new Foo(); foo.setField1(i); foo.setField2("foo" + i); foo.setField3(System.currentTimeMillis()); producer.newMessage().value(foo).send(); } producer.close(); pulsarClient.close(); } }代码要点:
PulsarClient.builder().serviceUrl("pulsar://localhost:6650"):连接本地 standalone 的 broker 二进制服务端口(6650);AvroSchema.of(Foo.class):基于 POJO 的 getter/setter 结构自动推导 Avro Schema 并注册到 Pulsar,Foo的三个字段field1(int)、field2(String)、field3(long)即会成为表的三列;- 循环发送 1000 条消息到
test_topic,Topic 位于public/default命名空间。
发送完成后,回到presto>提示符即可查询:
presto> show tables in pulsar."public/default"; presto> select field1, field2, field3 from pulsar."public/default".test_topic limit 10;需要注意:由于test_topic首次发送消息时(发送第一条带 Schema 的消息后)Pulsar 会自动创建该 Topic,show tables可能出现短暂延迟,稍候再查即可。
理解 Pulsar SQL 的工作原理
在独立复现上述流程之后,了解其架构有助于更好地配置与排障。Pulsar SQL 的整体介绍见 Pulsar SQL Overview,核心事实如下:
- Pulsar SQL = Presto 引擎 + Presto Pulsar connector。connector 使 Presto worker 能够把 Pulsar 的 Topic 当作关系表查询,这是整个能力的核心。
- 数据直接从 BookKeeper 读取,不经过 broker。Pulsar 采用两级分段(two-level segment based)架构,Topic 数据以 segment 形式存储在 Apache BookKeeper 中,每个 segment 在多个 BookKeeper 节点上冗余复制。connector 让 Presto worker 直接从 BookKeeper 并发读取,因此查询吞吐可以随 BookKeeper 节点数量水平扩展。
- 查询滞后性说明:由于 SQL worker 绕过 broker 直接读 BookKeeper,而 broker 不会主动推进 LAC(Last Add Confirmed),SQL 只能读到所有 bookie 已知的 LAC 之前的 entry。若要读到更新的数据,可以在
broker.conf中设置bookkeeperExplicitLacIntervalInMills让 broker 周期性显式写入 LAC(详见 sql-deployment-configurations 中的说明)。
Connector 关键配置
connector 的配置集中在conf/presto/catalog/pulsar.properties(仓库中真实文件为 conf/presto/catalog/pulsar.properties),文档与仓库中的核心参数如下:
| 配置项 | 默认值 | 说明 |
|---|---|---|
connector.name | pulsar | connector 名称,即在show catalogs中展示的 catalog 名 |
pulsar.web-service-url | http://localhost:8080 | Pulsar broker 的 Web 服务地址(注意:该配置在仓库中标注为DEPRECATED,同时存在pulsar.broker-service-url) |
pulsar.zookeeper-uri | localhost:2181 | ZooKeeper 集群地址 |
pulsar.max-entry-read-batch-size | 100 | 单次读取的最少 entry 数量(文档中写作pulsar.entry-read-batch-size,当前仓库中的配置键为pulsar.max-entry-read-batch-size,对应源码 PulsarConnectorConfig.java 中的entryReadBatchSize = 100) |
pulsar.target-num-splits | 4(文档)/2(当前仓库默认) | 每次查询默认使用的 split 数量,影响查询并行度 |
从源码 PulsarConnectorConfig.java 可以看到,connector 还支持更多可调项,例如pulsar.max-split-message-queue-size(默认 10000)、pulsar.max-split-entry-queue-size(默认 1000)、BookKeeper 客户端线程数、Managed Ledger 缓存大小(pulsar.managed-ledger-cache-size-MB,默认 0 即关闭)、TLS 与认证相关配置等。多 broker / 多 ZooKeeper 场景下,可以用逗号分隔多个地址,例如:
pulsar.web-service-url=http://localhost:8080,localhost:8081,localhost:8082 pulsar.zookeeper-uri=localhost1,localhost2:2181进一步阅读
- 如果你想在已有 Presto 集群中接入 Pulsar,或部署多节点 Pulsar SQL 集群(coordinator + worker 的完整配置示例,以及
./bin/pulsar sql-worker --help的 launcher 参数说明),参见 Pulsar SQL configuration and deployment。 - 关于 Pulsar SQL 的架构与性能设计(两级分段存储、BookKeeper 并发读取),参见 Pulsar SQL Overview。
- 关于 standalone 的安装、启动与停止细节,参见 Set up a standalone Pulsar locally。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar SQL 入门实战:从零开始用 Presto 查询 Pulsar 数据
Apache Pulsar SQL 入门实战:从零开始用 Presto 查询 Pulsar 数据 本指南面向 Apache Pulsar 2.2.1 及后续版本
消息队列后端流处理Apache Pulsar SQL 快速入门:用标准 SQL 查询 Topic 数据
Apache Pulsar SQL 快速入门:用标准 SQL 查询 Topic 数据 Apache Pulsar SQL 让开发者可以直接用标准 SQL 语句查
消息队列后端流处理Apache Pulsar SQL 入门实战指南:用 Presto 查询 Pulsar 中的消息数据
Apache Pulsar SQL 入门实战指南:用 Presto 查询 Pulsar 中的消息数据 导读 本指南基于 Pulsar SQL 的官方快速入门文档
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考