用SQL查MQ:Trino、PGMQ、KsqlDB、Proton、OpenLooKeng、KarelDB
概述用SQL查MQ算是鬼点子吧这个方向已形成一条相当成熟的技术路线思想大致分成三条路径流式SQL引擎面向实时数据流。最贴合SQL查询MQ思想核心思路把MQ的Topic当成一张可持续查询的表SQL不是一次性查完而是持续订阅、不断产出结果。代表ksqlDB、Apache Flink。联邦查询引擎用SQL跨源查存量消息。解决需求把MQ里堆积的历史消息用SQL查询常用于排查问题、追踪消息轨迹、做临时分析。代表Trino、Apache Pulsar SQL消息过滤/检索内的SQL能力不是外部引擎来查MQ而是MQ自身提供SQL风格的过滤语法属于轻量级的内置方案。代表RocketMQ、KarelDBRocketMQ在消费者订阅时支持MessageSelector.bySql()用类SQL92的表达式如TAGS IS NOT NULL AND age BETWEEN 18 AND 60在Broker端过滤消息属性相比传统Tag过滤粒度更细。实现原理在Broker里用表达式解析对消息属性做过滤需要开启enablePropertyFiltertrue。Trino官网Java开发、开源GitHub13.1K Star3.7K Fork分布式SQL查询引擎原PrestoSQL从Presto分叉而来本身是分布式MPP SQL查询引擎。核心价值联邦查询跨源查询能处理从GB到PB级的数据通过连接器Connector机制把Kafka、Pulsar、MySQL、Hive、MongoDB等统统抽象成Catalog里的表。配好kafka.properties后Kafka的每条消息在Trino里就是一行可直接SELECT * FROM kafka.default.user_events WHERE ...支持跨源JOIN不同数据源的数据。适合数据仓库、日志分析、数据湖查询、交互式分析、BI场景。核心特性无共享架构Shared NothingWorker完全独立无磁盘写依赖全内存流水线各Operator间通过内存Page直接传递数据动态代码生成运行时生成优化字节码消除虚函数调用特定数据类型特化异步I/O模型网络与计算重叠非阻塞数据获取流水线气泡最小化高速查询弹性线性扩展低成本使用采用无共享架构通过列式内存格式和动态代码生成优化性能并提供丰富的连接器实现计算存储分离(存算分离)最大化下推优化以提升效率。核心组件Coordinator协调器接收SQL请求解析生成分布式执行计划调度Task到Worker监控查询状态Worker工作节点执行Task数据扫描、过滤、聚合等操作通过Driver驱动多个Operator最小执行单元Connector连接器解耦计算与存储通过插件支持新数据源关键接口getSplits()数据分片、getPage()。查询执行流程从SQL到分布式计算StorageMetastoreWorkerCoordinatorClientStorageMetastoreWorkerCoordinatorClient提交 SQL 查询获取元数据生成执行计划分发查询任务读取数据分片本地计算处理返回部分结果结果聚合返回最终结果流程解读解析与逻辑计划生成SQL Parser将SQL文本解析成抽象语法树ASTLogical Planner结合元数据如表结构、统计信息将AST转换成逻辑执行计划物理计划与分布式调度Distributed Planner将逻辑计划拆分成多个Stage阶段构成一个树状依赖结构Stage Scheduler将Stage分发到各个Worker节点上执行任务执行与流水线处理Task Executor每个Stage在Worker上表现为一或多个TaskOperator PipelineTask内部由多个Operator操作符如Scan、Filter、Join、Aggregation组成数据在这些Operator之间以“流水线”的方式传递中间结果通常不落盘直接在内存中流转这是Trino速度快的重要原因。PGMQ项目主页即官方文档基于PG、完全使用SQL语句实现的轻量级开源GitHub5.1K Star141 ForkMQ工具可实现类似AWS SQS或RSMQRedis Simple Message Queue的接口和功能。功能特性轻量级不需要后台工作进程或外部组件依赖只需安装PG扩展插件支持PG 14可靠性能够确保在可见性超时之内精确传递一次消息到消费者API兼容AWS SQS、RSMQ消息队列APIFIFO遵循FIFO原则支持消息分组和顺序处理分区队列集成pg_partman插件按时间或ID自动分表存储消息冷数据可自动淘汰性能提升通过创建非日志表create_unlogged可极大地提高消息队列处理时间但是可能在系统崩溃时丢失队列数据显式删除消息可一直存储在队列中直到明确删除为止消息重放消息可被归档用于长期保留和重放客户端驱动官方提供Rust、Python驱动社区提供Dart、Go、Elixir、Java、Kotlin、JavaScript、.NET、Ruby、PHP等语言驱动实战基于Docker部署dockerrun-d--namepgmq-postgres-ePOSTGRES_PASSWORDpostgres-p5432:5432 ghcr.io/pgmq/pg18-pgmq:v1.10.0命令行使用示例// 创建队列SELECTpgmq.create(my_queue);// 发送消息SELECT*frompgmq.send(queue_namemy_queue,msg{foo: bar2},delay5);// 读取消息SELECT*FROMpgmq.read(queue_namemy_queue,vt30,qty2);// 弹出消息SELECT*FROMpgmq.pop(my_queue);// 归档消息SELECTpgmq.archive(queue_namemy_queue,msg_id2);// 删除消息SELECTpgmq.delete(my_queue,2);// 删除队列SELECTpgmq.drop_queue(my_queue);每个消息队列都对应pgmq模式下的一个表表的命名规则为q_前缀加上消息队列名。以上语句创建的表名为pgmq.q_my_queue。消息内容使用JSON格式查询返回的是消息编号。发送消息时delay可选指定延迟时间秒数或时间戳在此之前该消息对消费者不可见。读取消息时qty指定条数vt指定消息不可见时长。如果这些消息在vt秒内不被删除或归档将会重新可见且可被其他消费者读取。弹出(POP)表示消费者读取消息之后立即从队列中删除相应的消息。归档消息意味着从队列中删除消息并将其插入归档表中消息队列my_queue对应的归档表为pgmq.a_my_queue。KsqlDB官网前身是KSQL基于Kafka Streams API、几乎纯Java编写、开源GitHub314 Star1K Fork实时数据流处理引擎提供强大且易用的SQL交互方式来处理Kafka数据流而无需编写代码。具备高扩展、高弹性、容错等优良特性并支持多种流式处理操作如数据过滤、转化、聚合、连接、窗口化和会话化。提供强大的功能来处理实时数据流适用于各种需要高效数据处理的场景。概念Event事件通常被称为“行”类似于关系数据库中的一行。Stream流流是不可变的、仅可追加的历史数据集合。一旦将一行插入流中就无法更改。流中的每一行数据存储在特定的分区中并且每行隐式或显式地拥有一个代表其身份的键Table表表是可变的、分区的集合其内容会随时间而变化。表通过利用每一行的键来工作给定键的最后一行表示该键标识的最新信息。实战部署形态是独立的KSQL Server进程集群可动态扩容。基于Docker Compose部署version:2services:ksqldb-server:image:confluentinc/ksqldb-server:0.15.0hostname:ksqldb-servercontainer_name:ksqldb-serverports:-8088:8088environment:KSQL_LISTENERS:http://0.0.0.0:8088KSQL_BOOTSTRAP_SERVERS:192.168.1.87:9092KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE:trueKSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE:trueksqldb-cli:image:confluentinc/ksqldb-cli:0.15.0container_name:ksqldb-clidepends_on:-ksqldb-serverentrypoint:/bin/shtty:true启动Docker容器并连接ksqlDBdocker-composeup-ddockerexec-itksqldb-cli ksql http://ksqldb-server:8088基于Kafka的topic创建Stream并查询Stream中的数据CREATESTREAM cr7_topic_stream(orderAmountINTEGER,orderIdINTEGER,productIdINTEGER,productNumINTEGER)WITH(kafka_topiccr7-topic,value_formatjson);SELECT*FROMcr7_topic_stream EMIT CHANGES;创建Table并查询其中的数据CREATETABLEcr7_topic_table(orderAmountINTEGER,orderIdINTEGER,productIdINTEGER,productNumINTEGER,kafkaProducerKeyVARCHARPRIMARYKEY)WITH(kafka_topiccr7-topic,value_formatjson);SELECT*FROMcr7_topic_table EMIT CHANGES;ProtonTimeplus开源GitHub2.2K Star110 ForkKsqlDB替代方案用C编写、可对接多消息系统集成ClickHouse的分析能力。实战基于命令行安装curlhttps://install.timeplus.com|sh在ClickHouse里基于MergeTree引擎创建表CREATETABLEevents(_tp_time DateTime64(3),url String,method String,ip String)ENGINEMergeTree()PRIMARYKEY(_tp_time,url);https://zhuanlan.zhihu.com/p/687208908Apache Pulsar SQLPulsar的数据存在BookKeeper里本身没有表的概念但官方直接内置Pulsar SQL本质是打包一个Trino服务端并预装Pulsar插件提供sql-worker和sql两个命令。启动后就能对Topic直接执行SQL插件启动时通过Pulsar-Admin接口拉取Schema和分区元数据然后创建只读BookKeeper客户端按条件过滤数据实现像查表一样查消息。OpenLooKeng官网曾用名OpenHetu华为开源GitHub571 Star417 Fork使用Java开发、面向大数据库的高性能数据虚拟化引擎。提供统一SQL接口具备跨数据源/数据中心分析能力以及面向交互式、批、流等融合查询场景。同时增强前置调度、跨源索引、动态过滤、跨源协同、水平拓展等能力。基于SQL引擎Presto内置Kafka等连接器支持把Kafka Topic当表查。提供Coordinator AA高可靠、可扩展的数据源connector框架中文文档。核心优势跨数据中心数据分析统一SQL接口访问跨数据中心、跨云的数据源极简的跨源数据分析体验统一的SQL接口访问多种数据源易扩展数据源可通过增加Connector来增加数据源采集变连接、数据零搬迁实战命令行安装wget-O- https://download.openlookeng.io/install.sh|bashKarelDBJava开发、几乎完全构建在Kafka之上的开源GitHub389 Star27 Fork关系数据库MQ即存储方向的激进尝试。支撑组件Apache CalciteSQL引擎负责SQL解析、优化、执行Apache Omid事务管理MVVC支持Apache Avro序列化和模式演化schema evolutionApache Kafka持久化存储使用KCache嵌入式键值存储Apache AvaticaJDBC支持SQL类型boolean、integer、bigint、real、double、varbinary、varchar、decimal、date、time、timestamp实战Maven项目引入如下依赖dependencygroupIdio.kareldb/groupIdartifactIdkareldb-core/artifactIdversion1.0.0/version/dependency以内嵌模式工作PropertiespropertiesnewProperties();properties.put(schemaFactory,io.kareldb.schema.SchemaFactory);properties.put(parserFactory,org.apache.calcite.sql.parser.parserextension.ExtensionSqlParserImpl#FACTORY);properties.put(schema.kind,io.kareldb.kafka.KafkaSchema);properties.put(schema.kafkacache.bootstrap.servers,bootstrapServers);properties.put(schema.kafkacache.data.dir,/tmp);try(ConnectionconnDriverManager.getConnection(jdbc:kareldb:,properties);Statementsconn.createStatement()){s.execute(create table books (id int, name varchar, author varchar));s.executeUpdate(insert into books values(1, The Trial, Franz Kafka));ResultSetrss.executeQuery(select * from books);// 省略}