ARTICLE DETAIL

资讯详情

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

Apache Pulsar SQL 快速入门:用 Presto 引擎在 Pulsar 上执行 SQL 查询

Apache Pulsar SQL 快速入门:用 Presto 引擎在 Pulsar 上执行 SQL 查询 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文是 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.namepulsarsystem是 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 下有哪些表Topicpresto 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.wileygmail.com | genesisw | y9D2dtU3 | 959-197-1860 | 71 | genesis.wileyinterdemconsulting.eu | 880-58-9247 | Brayden | | Stanton | brayden.stantonyahoo.com | braydens | ZnjmhXik | 220-027-867 | 81 | brayden.stantonsupermemo.eu | 604-60-7069 | Benjamin | Julian | Velasquez | benjamin.velasquezyahoo.com | benjaminv | 8Bc7m3eb | 298-377-0062 | 21 | benjamin.velasquezhostesltd.biz | 213-32-5882 | Michael | Thomas | Donovan | donovanmail.com | michaeld | OqBm9MLs | 078-134-4685 | 55 | michael.donovanmemortech.eu | 443-30-3442 | Brooklyn | Avery | Roach | brooklynroachyahoo.com | broach | IxtBLafO | 387-786-2998 | 68 | brooklyn.roachwarst.biz | 085-88-3973 | Skylar | | Bradshaw | skylarbradshawyahoo.com | skylarb | p6eC6cKy | 210-872-608 | 96 | skylar.bradshawflyhigh.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_testTopicPulsar 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(); ProducerFoo 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 二进制服务端口6650AvroSchema.of(Foo.class)基于 POJO 的 getter/setter 结构自动推导 Avro Schema 并注册到 PulsarFoo的三个字段field1int、field2String、field3long即会成为表的三列循环发送 1000 条消息到test_topicTopic 位于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 会自动创建该 Topicshow 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 不会主动推进 LACLast Add ConfirmedSQL 只能读到所有 bookie 已知的 LAC 之前的 entry。若要读到更新的数据可以在broker.conf中设置bookkeeperExplicitLacIntervalInMills让 broker 周期性显式写入 LAC详见 sql-deployment-configurations 中的说明。Connector 关键配置connector 的配置集中在conf/presto/catalog/pulsar.properties仓库中真实文件为 conf/presto/catalog/pulsar.properties文档与仓库中的核心参数如下配置项默认值说明connector.namepulsarconnector 名称即在show catalogs中展示的 catalog 名pulsar.web-service-urlhttp://localhost:8080Pulsar broker 的 Web 服务地址注意该配置在仓库中标注为DEPRECATED同时存在pulsar.broker-service-urlpulsar.zookeeper-urilocalhost:2181ZooKeeper 集群地址pulsar.max-entry-read-batch-size100单次读取的最少 entry 数量文档中写作pulsar.entry-read-batch-size当前仓库中的配置键为pulsar.max-entry-read-batch-size对应源码 PulsarConnectorConfig.java 中的entryReadBatchSize 100pulsar.target-num-splits4文档/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-urlhttp://localhost:8080,localhost:8081,localhost:8082 pulsar.zookeeper-urilocalhost1,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。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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 的官方快速入门文档消息队列后端流处理上一篇Higress网关TLS安全加固实践协议版本与密码套件配置指南下一篇openEuler/QA版本测试流程揭秘从开发到发布的完整质量保障指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表