ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Cube ksqlDB 驱动(@cubejs-backend/ksql-driver)完全指南:流式预聚合、Kafka 直连与 SQL 方言适配

Cube ksqlDB 驱动(@cubejs-backend/ksql-driver)完全指南:流式预聚合、Kafka 直连与 SQL 方言适配 Cube ksqlDB 驱动cubejs-backend/ksql-driver完全指南流式预聚合、Kafka 直连与 SQL 方言适配【免费下载链接】cube Cube Core is open-source semantic layer for AI, BI and embedded analytics项目地址: https://gitcode.com/gh_mirrors/cu/cubeksqlDB 是构建在 Apache Kafka 之上的流处理数据库而cubejs-backend/ksql-driver是 Cube 开源语义层中专门对接 ksqlDB 的数据库驱动包。本文以该包 CHANGELOG 的演进记录为主线结合 KsqlDriver.ts、KsqlQuery.ts 源码与官方 ksqlDB 接入文档完整讲解该驱动的配置方式、架构原理、SQL 方言适配、流式预聚合与 Kafka 直连模式以及历史上关键能力与修复的来龙去脉。读完本文你将能够正确配置 ksqlDB 数据源、理解为什么 SELECT 只能走 Cube Store、掌握stream_offset/unique_key_columns/output_column_types等流式预聚合核心参数并搭建批流一体的 lambda 预聚合。驱动概况与支持状态cubejs-backend/ksql-driver是 Cube 众多数据库驱动之一其 package.json 中描述为 Cube ksql database driver要求 Node.js 20运行时依赖cubejs-backend/base-driver、cubejs-backend/schema-compiler、cubejs-backend/shared以及axios调用 ksqlDB REST API、kafkajs连接 Kafka broker和async-mutex串行化建删表操作。需要特别注意的是该包的社区支持属性README.md 明确标注community supported即社区维护、使用自担风险Cube 团队没有主动的功能开发计划但不影响其他模块的 bug 修复覆盖并长期招募维护者。这是评估生产环境选用该驱动时的重要前提。从 CHANGELOG 看该驱动自 0.28.472021-10-22ksql support首次引入后持续演进当前版本为 1.7.422026-09-18。其版本号与整个 Cube monorepo 保持同步绝大多数条目是 Version bump only真正属于 ksql 驱动自身的实质性变更集中在流式能力、预聚合、并发控制与查询方言修正几个方向。环境变量与连接配置驱动通过cubejs-backend/shared的getEnv读取环境变量在 KsqlDriver.ts 构造函数 中完成解析。官方 ksqldb.mdx 给出的基础.env配置如下CUBEJS_DB_TYPEksql CUBEJS_DB_URLhttps://xxxxxx-xxxxx.us-west4.gcp.confluent.cloud:443 CUBEJS_DB_USERusername CUBEJS_DB_PASSpassword各变量含义如下表依据 ksqldb.mdx 环境变量表环境变量说明取值必填CUBEJS_DB_URLksqlDB 主机 URL含端口合法数据库主机 URL✅CUBEJS_DB_USER连接用户名Confluent Cloud 下为 API key 名合法用户名✅CUBEJS_DB_PASS连接密码Confluent Cloud 下为 API secret合法密码✅CUBEJS_DB_KAFKA_HOSTKafka broker 地址用于 Kafka streams 模式多个 broker 用逗号分隔合法 Kafka broker URL❌CUBEJS_DB_KAFKA_USERKafka broker 认证用户名SASL PLAIN合法 Kafka 用户名❌CUBEJS_DB_KAFKA_PASSKafka broker 认证密码SASL PLAIN合法 Kafka 密码❌CUBEJS_DB_KAFKA_USE_SSL为true时启用 SASL_SSLtrue、false❌CUBEJS_CONCURRENCY对数据源的并发查询数合法数字❌源码中driverEnvVariables()KsqlDriver.ts声明的基础变量为CUBEJS_DB_URL、CUBEJS_DB_USER、CUBEJS_DB_PASS其余 Kafka 相关变量通过dbKafkaHost、dbKafkaUser、dbKafkaPass、dbKafkaUseSsl的getEnv调用读取且支持dataSource与preAggregations两种上下文这意味着 ksql 连接既可以作为具名数据源多数据源场景也可以配置独立的预聚合专用连接对应 1.6.34 的 pre-aggregation-specific data source configuration 能力。若使用 Confluent Cloud需要为 ksqlDB 集群与 Kafka 集群分别生成 API keyksqlDB 凭据填CUBEJS_DB_USER/CUBEJS_DB_PASSKafka 凭据填CUBEJS_DB_KAFKA_USER/CUBEJS_DB_KAFKA_PASS。架构REST API Kafka 双通道从 KsqlDriver 类定义 可以看到驱动内部维护了两条通道ksqlDB REST API 通道apiQuery()通过axios.post向${url}/ksql端点发送{ ksql: SQL; }并使用 HTTP Basic 认证auth: { username, password }。所有元数据操作SHOW TABLES、SHOW STREAMS、DESCRIBE和非查询语句CREATE TABLE ... AS SELECT、DROP TABLE都走这里。请求失败时会从响应中提取message/statusCode拼装错误信息。Kafka 直连通道可选当配置了kafkaHost时构造函数会用kafkajs创建Kafka客户端clientId: Cubebroker 列表由逗号分隔的kafkaHost拆分而来支持ssl与 SASL PLAIN 认证kafkaUser/kafkaPassword。该通道用于把流数据直接拉入 Cube Store绕过 ksqlDB 的 REST 流式接口。并发模型默认并发为 1getDefaultConcurrency()KsqlDriver.ts返回 1。这并非随意设定而是 ksqlDB 引擎的限制——0.33.242023-06-05的 Bug Fix 条目明确写道Reduce default concurrency to 1 as ksql doesnt support multiple queries at the same time right nowksqlDB 当前不支持同时执行多个查询。此前 0.30.302022-07-05还统一了各驱动的默认并发值并引入集中式并发设置centralized concurrency setting用户仍可通过CUBEJS_CONCURRENCY覆盖。连接校验与表操作testConnection()先执行SHOW VARIABLES若配置了 Kafka 客户端还会执行admin().connect()/disconnect()校验 broker 可达性对应 0.32.2 的 connection validation and logging。createSchemaIfNotExists()为空实现——ksqlDB 没有 schema 概念。dropTable()使用async-mutex串行化执行DROP TABLE \ DELETE TOPIC避免并发删除竞态对应 0.32.27 的 Drop table with the topic to avoid orphaned topics删除表时连带删除底层 topic防止遗留孤儿 topic。quoteIdentifier()与escapeColumnName()统一使用反引号包裹标识符ksqlDB 语法与 Presto 同源标识符引用用反引号而非双引号。tableDashName()将表名中的.替换为-以适配 ksqlDB 对象命名约定。为什么 SELECT 必须走 Cube Store驱动的query()方法KsqlDriver.ts对select开头的查询直接抛出异常Select queries for ksql allowed only from Cube Store. In order to query ksql create pre-aggregation first.这条约束正是 1.1.92024-12-08Bug Fix 的落地内容ksqlDB 的 SELECT 是持续流式查询pull/push 语义与批式 SQL 不同不适合作为语义层的即时查询目标。正确用法是先基于 ksqlDB 的 stream/table 构建预聚合把数据落到 Cube Store再由 Cube Store 响应 SELECT。只有当启用了 Kafka 直连下载配置kafkaHost时Cube Store 才能直接从 Kafka topic 拉取数据完成预聚合构建进而支撑查询。SQL 方言适配层KsqlQueryKsqlQuery.ts 继承 schema-compiler 的BaseQuery是 Cube 语义层 SQL 生成器与 ksqlDB 方言之间的适配器承担了几乎全部方言差异的翻译工作时间函数GRANULARITY_TO_INTERVAL表定义了从second到year共 9 种粒度的时间截断表达式统一使用FORMAT_TIMESTAMPPARSE_TIMESTAMP组合day: (date) FORMAT_TIMESTAMP(${date}, yyyy-MM-ddT00:00:00.000), week: (date) FORMAT_TIMESTAMP(PARSE_TIMESTAMP(FORMAT_TIMESTAMP(${date}, YYYY-ww), YYYY-ww), yyyy-MM-ddT00:00:00.000), // month / quarter / year 依此类推convertTz()生成CONVERT_TZ(field, UTC, timezone)timeStampParam()生成PARSE_TIMESTAMP(?, yyyy-MM-ddTHH:mm:ss.SSSX, UTC)即参数化时间戳统一按 ISO-8601 UTC 解析timeGroupedColumn()将维度按粒度截断后包装为PARSE_TIMESTAMP(...)。表达式与语句模板sqlTemplatessqlTemplates()是方言适配的核心KsqlQuery.ts其关键决策在注释中都有交代标识符引号quotes.identifiers ksqlDB 用反引号而非双引号字符串拼接ksqlDB 没有||运算符concat_strings模板改写为CONCAT(...)concatStringsSql()同样生成CONCAT(a, b)时间戳字面量从 Cube 传来的 ISO-8601 时间戳带末尾Z而 ksqlDB 对带Z的字符串做 CAST 会静默返回 NULL注释引用 confluentinc/ksql#9094因此timestamp_literal模板先剥掉T/Z再用PARSE_TIMESTAMP(..., yyyy-MM-dd HH:mm:ss.SSS, UTC)显式指定 UTC 解析LIKE 模式like_pattern模板用CONCAT(%, value, %)拼通配符配合KsqlFilter.likeIgnoreCase()生成column ILIKE CONCAT(%, ?, %)GROUP BYksqlDB 不支持按位置引用列的 GROUP BYGROUP BY 1因此group_by_exprs模板改为输出完整表达式groupByClause()也直接以维度表达式 join 生成 GROUP BY 子句禁用能力delete templates.functions.WIDTH_BUCKETksqlDB 无此函数对应 1.7.19 中cubesql 支持 WIDTH_BUCKET 下推是 SQL API 侧的事驱动方言本身不支持、delete templates.statements.unionksqlDB 无集合运算对应 1.7.28 Push UNION down 后驱动通过禁用模板让 UNION 留在后处理阶段。类型与常用函数castToString()CAST(x as varchar(255))unixTimestampSql()UNIX_TIMESTAMP()preAggregationLoadSql()生成CREATE TABLE \ WITH (KEY_FORMATJSON) AS即 ksqlDB 建物化表时显式指定 JSON key 格式。 预聚合相关行为 preAggregationStartEndQueries()当预聚合声明了 partitionGranularity 时强制要求 buildRangeStart 与 buildRangeEnd缺失即抛错——这是 ksql 流式分区预聚合的硬性约束对应 1.3.34 Fix pre-agg partitions creation for Ksql preAggregationReadOnly()判断预聚合是否为只读路径——originalSql 且 SQL 形如 SELECT * FROM table或 rollup 且维度中包含主键primaryKey preAggregationAllowUngroupingWithPrimaryKey() 返回 true当 rollup 维度含主键时允许去掉 GROUP BY从而使生成的查询变成无分组的 SELECT ... FROM ...这是只读流式路径的前提下文详述 extractTableFromSimpleSelectAsteriskQuery()用正则从 SELECT * FROM table 中提取表名供下载路径识别源表。 SQL 参数转义双写引号防注入 1.7.122026-07-27修复了包括 ksql 在内的多个驱动pinot/dremio/ksql/databricks/hive/jdbc的 SQL 参数转义问题。驱动通过 prepareQueryWithParams() 调用 shared 包的 formatAnsi 完成转义其规则由 params-escaping.test.ts 完整固化 引号双写值中的 会被写成 使值无法逃出字符串字面量。例如 WHERE status ? 传入 a OR 11 --生成 WHERE status a OR 11 -- 反斜杠是普通字符ksqlDB源自 Presto 语法不把 \ 视为转义符因此保留原样不以 \\ 双写 数组参数逐元素转义IN (?) 传入 [its, b] 生成 IN (its, b) 多占位符按序替换、无参数查询原样透传。 测试类 TestKsqlDriver 直接以 new KsqlDriver({ url: http://localhost:8088 }) 实例化并调用 prepareQueryWithParams是典型的纯单测验证方式可在 test/unit 中查阅。 流式预聚合与 Kafka streams 模式 官方文档明确指出ksqlDB 只支持流式预聚合streaming pre-aggregations。这一定位在源码中体现为 capabilities() 返回 { streamingSource: true }KsqlDriver.tsquery-orchestrator 侧会依据该 capability 决定是否删除源临时表见 PreAggregationLoader.ts 中 dropSourceTempTable !capabilities?.streamingSource 的判断。 默认模式经 ksqlDB REST API 流式 默认情况下Cube 通过 ksqlDB REST API 做元数据发现SHOW TABLES/SHOW STREAMS/DESCRIBE并在预聚合构建时通过 REST API 把数据流进 Cube Store构建过程中 Cube 可能向 ksqlDB 下发 CREATE TABLE ... AS SELECT 等语句创建非只读预聚合。 Kafka streams 模式直连 Kafka 当设置 CUBEJS_DB_KAFKA_HOST 后自动激活 Kafka streams 模式 Cube 仍用 ksqlDB REST API 发现 stream/table 及其 schemaDESCRIBE 对每个对象从 ksqlDB 元数据解析出底层 Kafka topic 名 Cube Store 直接连接 Kafka broker 消费对应 topic不再经 ksqlDB 转发数据 所有预聚合走只读刷新路径Cube 不会向 ksqlDB 下发任何 CREATE TABLE/CREATE STREAM。 该模式适合不希望 Cube 在 ksqlDB 中创建对象、需要更高吞吐、ksqlDB 权限受限、或希望 Cube Store 直接消费 Kafka topic 的场景。对应源码为 getStreamingTableData()KsqlDriver.ts有 kafkaHost 时 streamingSource 类型为 kafka凭据含 user/password/host/use_ssl且 streamingTable 会被替换为 describe.sourceDescription.topic即真实 topic 名否则类型为 ksql凭据为 REST API 的 user/password/url。这与 CHANGELOG 中 0.31.322022-12-28Direct kafka download support for ksql streams and tables、0.30.42022-05-20Download streaming select * from table originalSql pre-aggregations directly 一脉相承。 历史演进上0.31.162022-11-23为流式消费增加了 offset earliest、数据回放replays与按分区流式消费per partition streaming支持0.31.112022-11-02引入 Cube Store 密封分区能力0.35.812024-09-12正式落地 ksql 与 rollup 预聚合的组合。 流式预聚合的关键配置项 以下是 ksqldb.mdx 中定义的流式预聚合核心参数 参数说明 read_only: trueCube 不在 ksqlDB 中创建任何对象数据直接从底层 Kafka topic 消费 stream_offset控制 Cube Store 从 Kafka topic 的何处开始消费latest 只消费建预聚合后新到达的消息earliest 从头回放整个 topic默认 latest。后续刷新时 Cube Store 自动从上次已处理 offset 继续不受该参数影响 unique_key_columns唯一标识一条记录的列名字符串非 member 引用配合内部序列列 __seq 构成表排序键用于读取与压缩时的去重——同一唯一键只保留序列号最高最新的一行 output_column_types声明 Cube Store 输出表的列类型Kafka streams 模式下配合 unique_key_columns 使用否则列名对不上会导致构建失败 partition_granularity build_range_start/build_range_end分区预聚合必需ksqlDB 不支持 NOW()需用 CURRENT_TIMESTAMP /- INTERVAL 表达式 只读流式路径的前置条件 要真正走通 Kafka streams 只读路径生成的 SQL 不能含 GROUP BYCube Store 的流后处理引擎不支持聚合。Cube 在 rollup 维度包含主键时会自动省略 GROUP BY——这正是 preAggregationReadOnly() 与 preAggregationAllowUngroupingWithPrimaryKey() 的配合逻辑。因此 流式预聚合的 dimensions 必须包含 cube 的全部主键列否则无法识别为无分组查询流式路径失效 cube 的 sql或 sql_table应引用已存在的 ksqlDB stream/tableschema 由 Cube 自动发现。 Topic 名匹配的已知限制 Kafka streams 模式下Cube Store 会解析生成的 select_statement 并拿 FROM 表名与真实 Kafka topic 名比对。在 Confluent Cloud 等托管平台上ksqlDB 对象名与 topic 名常常不一致例如 stream 名为 ORDER_EVENTS_STREAMtopic 却是 pksqlc-abc123ORDER_EVENTS_STREAM。当前驱动不会把 FROM 子句改写为解析出的 topic 名因此要求 ksqlDB 对象名与 Kafka topic 名完全一致含大小写否则构建报 Topic table ... is not found。创建 stream 时应显式指定同名 topic例如 CREATE STREAM ORDER_EVENTS_STREAM (...) WITH (KAFKA_TOPICORDER_EVENTS_STREAM, VALUE_FORMATJSON, ...); 这是 Kafka streams 模式的已知限制当 ksqlDB 使用默认 topic 命名策略时对象名与 topic 名一致不会触发该问题。 消息格式与时间戳处理 Cube Store 期望 Kafka 消息的 value 是 JSON 对象字段名与 cube 中维度/度量的 sql 列名大小写一致缺省字段为 null。消息 key 可选当 value 以 { 开头时key 会被解析为 JSON 对象作为唯一键列取值的兜底来源。 时间维度type: time接受两类值 字符串按 ISO 8601 / RFC 3339 解析支持如 2025-01-15T10:30:00.000Z、2025-01-15 10:30:00.000 UTC、2025-01-15 等格式 数字按epoch 毫秒解释非秒、非微秒例如 1736939400000 表示 2025-01-15T10:30:00.000Z。 非标准格式可用 PARSE_TIMESTAMP 在 cube 的 sql 中转换。若源列是 bigint 时间戳需注意单位换算CAST(x AS TIMESTAMP(6)) 把数字当微秒所以毫秒值要先乘 1000秒值乘 1000000微秒值直接 CAST。官方文档还提示调试 bigint 时间戳问题时可临时去掉 partition_granularity 跳过分区构建以简化验证。 流上过滤 流式 cube 的 sql 若带 WHERECube Store 会在每个微批次上直接应用投影与过滤在 Cube Store 内部执行不需要 ksqlDB 处理。但 SELECT 语句形态受限只接受 Projection Filter TableScan 形状Filter 可省略支持列引用、SELECT *、别名、比较运算符、AND/OR/NOT、IS [NOT] NULL、IN、BETWEEN、CASE WHEN、CAST、EXTRACT、SUBSTRING、标量函数、CONVERT_TZ、PARSE_TIMESTAMP/FORMAT_TIMESTAMP、date_trunc不支持 JOIN、子查询、GROUP BY/HAVING、聚合函数、ORDER BY、LIMIT/OFFSET、集合运算、窗口函数、CTE 等。SELECT 列表中非简单列引用的表达式必须有显式别名。 实战批流一体的 lambda 预聚合 ksqlDB 的典型定位是作为主数据仓库之外的附加实时数据源。官方文档给出的推荐模式是批量 cube读数仓历史数据按天增量建分区 流式 cubedata_source: ksql指向已存在的 ksqlDB stream用只读流式预聚合直连 Kafka topic lambda 预聚合合并两者。要点 用带装饰前缀的环境变量声明具名数据源例如 CUBEJS_DS_KSQL_DB_TYPEksql、CUBEJS_DS_KSQL_DB_KAFKA_HOST... 批量 cube 的预聚合用 rollup_lambdarollups 指向批量 rollup 与流式 rollup 流式 cube 的预聚合设置 read_only: true、stream_offset: latest、unique_key_columns含时间维度时写作 时间维度名_粒度如 created_at_second与 output_column_types build_range_start/build_range_end 用 CURRENT_TIMESTAMP /- INTERVAL 而非 NOW()。 完整 YAML 与 JavaScript 双版本示例见 ksqldb.mdx 数据建模章节。关于 unique_key_columns 与 indexes 的命名差异需要留意indexes 的 columns 用生成语句中的全限定别名cube_name__time_dimension_granularity而 unique_key_columns 只写 time_dimension_granularity且预聚合定义 indexes 时索引引用的每个维度都必须出现在 unique_key_columns 中否则构建会失败。 演进时间线驱动能力成长史 汇总 CHANGELOG 中与 ksql 驱动直接相关的实质性变更剔除纯版本同步条目 版本日期类型内容 1.7.422026-09-18同步当前最新版 1.7.372026-09-10Features迁移 TypeScript 6.0.3所有驱动支持命名 ESM 导出 1.7.282026-08-26FeaturescubesqlUNION 下推至数据源 1.7.192026-08-12FeaturescubesqlWIDTH_BUCKET SQL 下推支持 1.7.122026-07-27Bug Fixes修正 SQL 参数转义pinot/dremio/ksql/databricks/hive/jdbc 驱动 1.7.9 / 1.7.112026-07-23/26Featurescubesql支持解析仅含日期的 timestamp 字符串 1.6.342026-04-14Features支持预聚合专用数据源配置 1.3.342025-07-04Bug Fixes修复 Ksql 预聚合分区创建 1.1.92024-12-08Bug FixesSELECT 仅允许从 Cube Store 查询未启用 Kafka 下载时需先建预聚合 1.1.82024-12-05Bug Fixes修复 Kafka broker 列表处理 0.35.812024-09-12Featuresksql 与 rollup 预聚合 0.33.242023-06-05Bug Fixes默认并发降为 1ksql 不支持并发查询 0.32.272023-04-14Bug Fixes删表时连带删除 topic避免孤儿 topic 0.32.22023-03-07Features连接校验与日志 0.31.322022-12-28FeaturesCube Store 直接下载 ksql streams/tablesKafka 直连 0.31.162022-11-23Features支持 offset earliest、回放与按分区流式 0.31.112022-11-02FeaturesCube Store 密封分区 0.31.02022-10-03Features多数据源支持 0.30.42022-05-20Features直接下载 SELECT * FROM table 的 originalSql 预聚合 0.30.302022-07-05Fix/Feat统一各驱动默认并发集中式并发设置 0.29.512022-04-22Features查询语言支持 startsWith/endsWith 过滤器 0.28.482021-10-22Bug Fixes修复未加引号的 DESCRIBE 0.28.472021-10-22Featuresksql 支持首次引入 从这张时间线可以清晰看到驱动的演进主线先是 2021 年底完成基础对接元数据发现、方言适配2022 年补齐流式能力分区流式、offset 控制、Kafka 直连下载、多数据源、并发模型2023 年修复生命周期问题孤儿 topic、并发限制2024 年确立只读流式预聚合 Cube Store 查询的正确使用范式2025–2026 年则集中在转义安全、SQL 下推与工程基建TypeScript 6、ESM 导出上。 深入阅读 驱动实现KsqlDriver.tsREST/Kafka 双通道、预聚合下载、能力声明 方言适配KsqlQuery.ts时间函数、SQL 模板、只读预聚合判定 转义测试test/unit/params-escaping.test.ts 变更记录CHANGELOG.md 官方接入文档docs-mintlify/admin/connect-to-data/data-sources/ksqldb.mdx环境变量、Kafka streams 模式、lambda 预聚合完整示例 流式能力在上层编排器的消费点PreAggregationLoader.tsstreamOffset 参与版本键与刷新、streamingSource capability 判定 【免费下载链接】cube Cube Core is open-source semantic layer for AI, BI and embedded analytics 项目地址: https://gitcode.com/gh_mirrors/cu/cube创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表