StarRocks Elasticsearch Catalog 使用指南:免数据迁移直连 ES 的联邦查询方案
2026/9/16 18:27:57 网站建设 项目流程

StarRocks Elasticsearch Catalog 使用指南:免数据迁移直连 ES 的联邦查询方案

【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks

Elasticsearch Catalog 是 StarRocks 自 v3.1 起提供的免迁移联邦查询能力:无需将 Elasticsearch 中的索引数据导入 StarRocks,即可直接在 StarRocks 上用 SQL 对 ES 集群的索引进行多维分析与全文检索,并借助谓词下推把过滤逻辑下沉到 ES 执行。本文将从建目录语法、全部连接参数、谓词下推机制到esquery()高级查询,结合本仓库 FE 端源码(connector/elasticsearch)逐层拆解其实现原理与实战配置。

为什么需要 Elasticsearch Catalog

StarRocks 与 Elasticsearch 是两款侧重点不同的分析系统:StarRocks 擅长大规模分布式计算,可通过外表方式查询 ES 中的数据;Elasticsearch 则以全文检索能力见长。两者结合可以构成更完整的 OLAP 解决方案。借助 Elasticsearch Catalog,你可以在 StarRocks 上直接用 SQL 分析 ES 集群中所有索引数据,全程无需迁移数据。

一个值得注意的模型差异是:与其他数据源的 Catalog 不同,Elasticsearch Catalog 创建后只含有一个名为default_db的数据库,ES 中的每个索引(index)会自动映射为一张数据表并挂载到该数据库下。这一映射逻辑在源码 EsRestClient.listTables() 中实现:它会通过_cat/indices枚举索引、通过_aliases枚举别名,最终返回"索引 + 别名"的并集作为表集合,并过滤掉以.开头的系统索引(如.kibana_1.opendistro_security)。

创建 Elasticsearch Catalog

语法

CREATE EXTERNAL CATALOG <catalog_name> [COMMENT <comment>] PROPERTIES ("key"="value", ...)

参数说明

参数必填默认值说明
hostsES 集群连接地址,可指定一个或多个地址。StarRocks 会从该地址解析 ES 版本和索引分片分布。StarRocks 依据GET /_nodes/httpAPI 返回的地址与 ES 集群通信,因此hosts的取值必须与GET /_nodes/http返回的地址一致,否则 BE 或 CN 可能无法与 ES 集群通信。
type数据源类型,创建 ES Catalog 时必须设置为es
user开启 HTTP 基础认证时用于登录 ES 集群的用户名。需确保该用户有访问/cluster/state/nodes/http等路径及读取索引的权限。
password登录 ES 集群的密码。
es.type_doc索引类型。查询 ES 8 及以后版本的数据时无需配置(ES 8+ 已移除 mapping types)。
es.nodes.wan.onlyFALSE是否仅使用hosts指定的地址访问 ES 集群并拉取数据。true:不做数据节点嗅探,当 StarRocks 无法访问 ES 集群内数据节点地址时须设为truefalse:StarRocks 以hosts指定地址为入口嗅探索引分片所在的数据节点,生成执行计划后由 BE/CN 直连集群内数据节点拉取分片数据,网络互通时建议保留默认值false
es.net.sslFALSE是否允许使用 HTTPS 协议访问 ES 集群(仅 StarRocks v2.4 及以后支持)。true:HTTPS 与 HTTP 均可访问;false:仅支持 HTTP。
enable_docvalue_scanTRUE是否从 ES 列式存储(doc_values)中获取目标字段的值。多数场景下列式读取性能优于行式存储读取。
enable_keyword_sniffTRUE是否基于 ES 的 KEYWORD 类型字段嗅探 TEXT 类型字段。设为false时,StarRocks 在分词后进行匹配。

示例

CREATE EXTERNAL CATALOG es_test COMMENT 'test123' PROPERTIES ( "type" = "es", "es.type" = "_doc", "hosts" = "https://xxx:9200", "es.net.ssl" = "true", "user" = "admin", "password" = "xxx", "es.nodes.wan.only" = "true" );

创建成功后,可通过SHOW CATALOGS查看目录、SHOW DATABASES FROM es_test查看default_db数据库、SHOW TABLES FROM es_test.default_db查看自动映射出的索引表,随后即可像查询普通表一样执行SELECT分析。

参数背后的源码实现

上述参数并非只是文档中的说明文字,它们都有对应的 FE 端实现,理解实现能帮你更好地判断该如何取值。

  • 配置解析:参数解析集中在 EsConfig.java 中,通过@Config注解绑定hostsuserpasswordes.net.ssles.nodes.wan.onlyenable_docvalue_scanenable_keyword_sniff等键(键名常量定义于 EsTable.java)。连接器初始化时由 ElasticsearchConnector.bindConfig() 把配置转交给EsRestClient构造 HTTP 客户端。
  • hosts与节点容错EsRestClient在发起请求前会对节点地址做trim,并自动补全缺失的http://前缀(IPv6 请使用[addr]:port格式)。若某节点请求失败,会按nodes.length次循环切换到下一个节点重试,见 execute()。
  • 版本嗅探与兼容:StarRocks 通过请求 ES 根路径/读取响应中的version.number并解析为主版本号(EsMajorVersion.parse()),支持 0.x ~ 8.x 的识别,用于后续请求路径与查询 DSL 的兼容性判断。这正对应"StarRocks 可以从hosts地址解析 ES 版本"这一文档描述。
  • es.net.ssl:开启后使用信任所有证书的 TLS SocketFactory 构造 HTTPS 客户端(getOrCreateSSLClient())。
  • enable_docvalue_scan:默认true。源码注释(EsTable.java)引用 Solr 的 benchmark 指出:当返回字段较少时 DocValues 列式读取性能优于 stored_fields,但字段数量增多后差距缩小甚至反转,因此实现中还有一个内部上限max_docvalue_fields(默认 20,可通过max_docvalue_fields配置),超过该字段数会降级为从 stored_fields 提取,text类型字段默认禁用 doc_values 扫描。
  • 索引映射与分片发现EsRestClient通过GET /_nodes/http发现集群 HTTP 节点(getHttpNodes()),通过GET {index}/_search_shards获取索引分片所在节点(searchShards())。这就是es.nodes.wan.only=false时 BE/CN 能直连数据节点读取分片数据的底层依据;同时它也会通过_cat/indices/{index}?h=docs.count估算行数,供 StarRocks 的 CBO 优化器使用。

谓词下推(Predicate Pushdown)

StarRocks 支持将针对 ES 表的查询谓词下推到 ES 执行,缩小查询引擎与存储源之间的距离,从而提升查询性能。下推的 SQL 语法与生成的 ES 查询 DSL 对应关系如下:

SQL 语法Elasticsearch 语法
=term query
interms query
>=, <=, >, <range
andbool.filter
orbool.should
notbool.must_not
not inbool.must_not + terms
esqueryES Query DSL

在 FE 端,这一映射由 QueryConverter.java 完成。它以 AST 访问器的方式遍历查询表达式树,将=转为termQuery(visitBinaryPredicate())、IN转为termsQuery(visitInPredicate())、比较运算转为rangeQueryAND/OR/NOT分别转为bool.must / bool.should / bool.mustNot(visitCompoundPredicate())。除此之外,实现还额外支持了IS NULL(映射为 exists query 及其取反)和LIKE(把_%通配符转换为 ES wildcard 的?*)等下推规则。

值得注意的下推边界:下推过程会把"能下推的表达式"收集为remoteConjuncts(由 ES 处理),把"无法下推的表达式"收集为localConjuncts(仍由 StarRocks BE 处理),见 convert()。也就是说,谓词下推是"尽力而为"的优化,无法转换的过滤条件会自动留在 StarRocks 侧执行,不会导致查询失败。

使用 esquery() 下发 ES 原生查询

esquery()函数用于把 SQL 无法表达的 ES 查询(如 match、geo_shape 等)原样下推到 ES 进行过滤处理。其第一个参数是用于关联索引的列名,第二个参数是花括号({})包裹的、基于 ES Query DSL 的 JSON 表示;该 JSON只能且必须有且仅有一个根键,例如matchgeo_shapebool。在 FE 实现中,esquery(col, '{...}')会被转换为一个 RawQueryBuilder,直接将第二个参数的 JSON 原文拼入下推的查询 DSL(visitFunctionCall())。

  • Match 查询
SELECT * FROM es_table WHERE esquery(k4, '{ "match": { "k4": "StarRocks on elasticsearch" } }');
  • Geoshape 查询
SELECT * FROM es_table WHERE esquery(k4, '{ "geo_shape": { "location": { "shape": { "type": "envelope", "coordinates": [ [ 13, 53 ], [ 14, 52 ] ] }, "relation": "within" } } }');
  • Boolean 查询
SELECT * FROM es_table WHERE esquery(k4, ' { "bool": { "must": [ { "terms": { "k1": [ 11, 12 ] } }, { "terms": { "k2": [ 100 ] } } ] } }');

配合上一节的谓词下推表可以看到,esquery()是谓词下推的"白名单外"通道:凡是标准 SQL 下推覆盖不到的 ES 高级检索(全文 match、地理位置、嵌套 bool 组合等),都可以通过它直接以原生 DSL 形式下发,让 ES 的全文检索能力真正融入 StarRocks 的 SQL 分析链路。

使用说明与注意事项

  • 版本兼容:自 ES v5.x 起扫描数据的底层方式发生了变化,StarRocks 仅支持查询 Elasticsearch v5.x 及之后版本的数据。FE 端 EsMajorVersion 的版本解析范围(0.x ~ 8.x)与此约束一致。
  • 认证要求:StarRocks 仅支持查询开启了 HTTP 基础认证(HTTP basic authentication)的 ES 集群,请确保在 PROPERTIES 中正确配置userpassword,且该用户具备访问/cluster/state/nodes/http等元数据路径和读取索引的权限。
  • 性能提醒:某些查询(如包含count()的查询)在 StarRocks 上的执行速度会明显慢于直接在 ES 上执行,因为 ES 可以直接读取满足查询条件的文档数量元数据,而无需过滤实际数据。对这类聚合统计场景,可评估直接在 ES 侧完成或采用定期物化的方式。
  • 连接模式选择:若 BE/CN 与 ES 集群处于同一内网,建议保持es.nodes.wan.only = false,让 BE/CN 直连分片所在数据节点以获得更好性能;若存在网络隔离、无法访问 ES 内部数据节点地址,则须设置为true,此时所有流量都经由hosts指定的地址转发。

参考源码路径

  • 连接器入口与配置绑定:ElasticsearchConnector.java、EsConfig.java
  • REST 客户端(版本嗅探、节点发现、分片发现、索引列表):EsRestClient.java、EsMajorVersion.java
  • 谓词下推与 esquery 转换:QueryConverter.java、QueryBuilders.java
  • 表属性常量与 doc_values 策略:EsTable.java
  • 单元测试:ElasticsearchMetadataTest.java

【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询