wren-semantic-core 深度解析:基于 Apache DataFusion 的 MDL 语义查询引擎
【免费下载链接】WrenAIGenBI (Generative BI) for AI agents, an open-source, governed text-to-SQL through an open context layer that turns natural-language questions into trusted dashboards, charts, and SQL across 20+ data sources, such as BigQuery, Snowflake, PostgreSQL, ClickHouse, Amazon Redshift, Databricks and more.项目地址: https://gitcode.com/GitHub_Trending/wr/WrenAI
本文以 core/wren-core/core/README.md 为主线,结合
wren-semantic-core的源码、Cargo 配置、示例与测试,讲解如何将 MDL(Modeling Definition Language)清单与原始 SQL 结合,完成语义重写、访问控制与逻辑计划优化。读者将掌握该引擎的核心 API(AnalyzedWrenMDL、transform_sql等)、MDL 清单结构、三种执行模式以及 RLAC/CLAC 访问控制的底层实现。
在 Wren AI 的开源生态中,wren-semantic-core是位于 Wren AI 底层的 Rust 语义引擎——它面向 MCP 客户端与 AI Agent 提供"语义层(Semantic Layer)"能力。引擎接收一份MDL(Modeling Definition Language)清单与一条 SQL 查询,将查询经由语义层重写:解析模型、关系、指标与视图,应用行级/列级访问控制,最终产出一条优化后的逻辑计划。整个引擎构建在Apache DataFusion之上,既可用于在本地运行时直接执行查询,也可用于把语义化 SQL 反解析(unparse)为目标数据源可执行的 SQL。
Wren AI 整体架构图
从 misc/wren-ai-architecture.png 可以看到,MDL 语义建模(Models、Relationships、Calculations、Views)位于 Wren AI 的开放上下文层(Open Context Layer)内部,是连接上层 AI Agent 应用与下层 20+ 数据源的关键枢纽,而wren-semantic-core正是这一层的 Rust 实现。
一、包名与库名:wren-semantic-corevswren_core
理解这个项目的第一件事是区分发布名与导入名:
- 发布到 crates.io 的包名是
wren-semantic-core; - 而在 Rust 代码中导入的是
wren_core。
这一点在 core/wren-core/core/Cargo.toml 中有明确对应:[package] name = "wren-semantic-core",而[lib] name = "wren_core"。包描述为"Wren semantic engine — MDL-based semantic SQL layer and query planner built on Apache DataFusion",关键字为sql、semantic-layer、datafusion、mdl、query。
在仓库的成员工作区中,除了该核心 crate,还有基于它封装的 wren-core-py(Python 绑定)、wren-core-wasm(WebAssembly 构建)以及基础库 wren-core-base,后者提供了 MDL 的Manifest、builder、DataSource等底层定义,被wren_core直接 re-export。
二、安装与最小使用
在Cargo.toml中加入依赖:
[dependencies] wren-semantic-core = "0.1"最小用法(来自 core/wren-core/core/README.md):
use wren_core::mdl::AnalyzedWrenMDL; // Build an AnalyzedWrenMDL from a manifest, then transform SQL through the // semantic layer. See the API docs for the full flow.AnalyzedWrenMDL是语义分析后的核心数据结构。从 src/mdl/mod.rs 可以看到它的定义:
pub struct AnalyzedWrenMDL { pub wren_mdl: Arc<WrenMDL>, pub lineage: Arc<lineage::Lineage>, }其中WrenMDL持有解析后的Manifest、列引用映射(qualified_references)、已注册的物理表(register_tables)以及catalog_schema_prefix;Lineage(见 src/mdl/lineage.rs)则通过 petgraph 有向图维护计算字段的血缘依赖(source_columns_map、required_fields_map、required_dataset_topo),用于决定查询某列时实际需要扫描哪些源列、需要展开哪些关系链。
Feature 开关:multi-thread
Cargo.toml 中默认启用multi-threadfeature:
[features] default = ["multi-thread"] # Enable multi-threaded tokio runtime (required for sync transform_sql). # Disable for WASM builds where only single-threaded async is available. multi-thread = ["tokio/rt-multi-thread"]- 启用时,
transform_sql这个同步包装函数可用(内部自建 tokio runtime 并block_on); - 在 WASM 构建(如 wren-core-wasm)下必须关闭,此时只能使用
transform_sql_with_ctx异步接口。
transform_sql的#[cfg(feature = "multi-thread")]属性在 src/mdl/mod.rs 中显式标出,并注释"Not available on WASM — usetransform_sql_with_ctxdirectly in async context."
三、MDL 清单:语义层的输入模型
wren-semantic-core不直接理解任意数据库表,而是以 MDL 清单(Manifest)为输入。下面是一份来自测试数据的完整清单 core/wren-core/core/tests/data/mdl.json(节选),它定义了catalog、schema、三个 model、两个 relationship、一个 view 以及dataSource:
{ "catalog": "test", "schema": "test", "models": [ { "name": "customer", "tableReference": { "catalog": "", "schema": "", "table": "customer" }, "columns": [ { "name": "c_custkey", "type": "integer" }, { "name": "c_name", "type": "varchar" }, { "name": "custkey_plus", "type": "integer", "expression": "c_custkey + 1", "isCalculated": true }, { "name": "orders", "type": "orders", "relationship": "CustomerOrders" } ], "primaryKey": "c_custkey" } ], "relationships": [ { "name": "CustomerOrders", "models": ["customer", "orders"], "joinType": "ONE_TO_MANY", "condition": "customer.c_custkey = orders.o_custkey" } ], "views": [ { "name": "customer_view", "statement": "select * from test.test.customer" } ], "dataSource": "mysql" }关键要素与引擎的对应关系:
| MDL 要素 | 含义 | 引擎中的处理 |
|---|---|---|
models[].tableReference | model 对应的物理表 | 构建WrenDataSource并注册为 DataFusion TableProvider(src/mdl/mod.rs) |
columns[].isCalculated+expression | 计算列/指标 | 血缘分析记录其源列,ModelAnalyzeRule重写为底层表达式(src/mdl/lineage.rs) |
columns[].relationship | 关系列 | 参与关系链解析与 join 生成 |
relationships[] | 模型间 join 条件与类型 | RelationChain计算关系链拓扑 |
views[] | 语义视图 | ExpandWrenViewRule首先展开视图 |
dataSource | 目标数据源 | 决定使用哪个方言的 UDF 集合与WrenDialect |
从源码看,infer_and_register_remote_table(src/mdl/mod.rs)会为每个ModelSource::TableReference的 model 推断 schema 并注册WrenDataSource,为ModelSource::RefSql的 model 从物理列构建 schema。值得注意的细节是:CLAC 校验失败(无权访问)的列不会进入注册的 schema——这正是后面"列级访问控制"得以生效的机制:无权列对查询计划而言"根本不存在",从而报出列不存在错误。
另外,src/mdl/mod.rs 的analyze_with_url_tables提供了一个面向 WASM/文件数据源的变体:对于LocalFile、MinioFile、S3File、GcsFile等文件型数据源,跳过 WebDAV PROPFIND 列目录操作,直接把tableReference解析为 URL,按扩展名假定 Parquet 格式,仅通过infer_schema(GET + Range 读 Parquet footer)即可注册表——遵循 DuckDB-WASM 的"已知 URL + 已知格式 = 无需列目录"模式。
四、查询重写管线:从 SQL 到目标 SQL
引擎对外最核心的入口有两个(src/mdl/mod.rs):
// 同步入口(需要 multi-thread feature) pub fn transform_sql( analyzed_mdl: Arc<AnalyzedWrenMDL>, remote_functions: &[RemoteFunction], properties: HashMap<String, Option<String>>, sql: &str, ) -> Result<String> // 异步入口(可配合自定义 SessionContext) pub async fn transform_sql_with_ctx( ctx: &SessionContext, analyzed_mdl: Arc<AnalyzedWrenMDL>, remote_functions: &[RemoteFunction], properties: SessionPropertiesRef, sql: &str, ) -> Result<String>其中properties的类型是Arc<HashMap<String, Option<String>>>,即会话属性(如x-wren-timezone),会被传给访问控制条件中的@name占位符。remote_functions是用户自定义的远程函数(UDF/UDAF/UDWF),引擎在重写前会将其注册进 SessionContext(src/mdl/mod.rs)。
transform_sql_with_ctx的完整调用链(从源码逐行还原):
- 注册远程函数:按
FunctionType::Scalar | Aggregate | Window分别注册ByPassScalarUDF、ByPassAggregateUDF、ByPassWindowFunction;注意 DataFusion 解析时会小写化函数名,因此注册时会同时保留原名与转小写的别名,以保证 SQL 生成时输出原名; apply_wren_on_ctx(src/mdl/context.rs):派生一个新的SessionContext——设置default_null_ordering = nulls_last、关闭标识符规范化、把默认 catalog/schema 指向 MDL 的catalog/schema、开启 information_schema,并按x-wren-timezone覆盖时区;随后按Mode::Unparse挂载 analyzer/optimizer 规则并注册 MDL 表;ctx.state().create_logical_plan(sql):把原始 SQL 解析成 DataFusion 逻辑计划;- 失败兜底:若建计划失败,调用
permission_analyze(src/mdl/mod.rs)以Mode::PermissionAnalyze重新分析,区分"真的是权限错误"(返回友好的WrenError)与"普通 SQL 错误"(返回原始错误); ctx.state().optimize(&plan):执行优化规则;- 用
WrenDialect::new(&data_source)构造针对目标数据源(MySQL、BigQuery、PostgreSQL……)的方言,配合Unparser把逻辑计划反解析回 SQL,最后去掉 MDL 自带的catalog.schema.前缀(catalog_schema_prefix替换)。
三种执行模式(Mode)
src/mdl/context.rs 定义了三种模式,它们决定挂载哪些规则:
| 模式 | 用途 | 特点 |
|---|---|---|
LocalRuntime | 本地执行 | 由 DataFusion 直接执行查询,使用 DataFusion 原生TypeCoercion,无自定义 optimizer 规则 |
Unparse | 生成 SQL | 面向其他 SQL 引擎输出语句,使用WrenTypeCoercion,并应用一组为 unparser 裁剪过的 optimizer 规则 |
PermissionAnalyze | 权限诊断 | 仅在 Unparse 报错时触发,用于判断错误是否由权限拒绝引起 |
Analyzer 与 Optimizer 规则编排
Unparse 模式的 analyzer 规则顺序为(src/mdl/context.rs):
ExpandWrenViewRule——必须先执行,把视图展开为基础模型;ModelAnalyzeRule——解析 model 扫描、计算列、关系列,重写逻辑计划(对应 src/logical_plan/analyze/model_anlayze.rs);ModelGenerationRule——按血缘生成必要的 join 关系;TimestampSimplify——简化时间戳比较,且必须置于类型强制转换之前,以便简化结果按需转型;WrenTypeCoercion——语义层专用的类型强制转换。
Optimizer 规则(src/mdl/context.rs)则从 DataFusion 默认规则中精心挑选并显式禁用了一部分,原因都写在注释里:例如禁用EliminateNestedUnion(unparser 只支持两路 union)、禁用PushDownFilter(避免 BigQuery 的 datetime/timestamp cast 被移除)、禁用EliminateCrossJoin(避免生成无 join 条件的非法 SQL)等。最终启用:EliminateJoin、ExtractEquijoinPredicate、EliminateDuplicatedExpr、EliminateFilter、PropagateEmptyRelation、FilterNullJoinKeys、EliminateOuterJoin、EliminateGroupByConstant。
这些规则的组织在 src/logical_plan/analyze/mod.rs 与 src/logical_plan/optimize/mod.rs 中可见,分别包含access_control、expand_view、model_anlayze、model_generation、plan、relation_chain、scope与simplify_timestamp、type_coercion模块。
五、访问控制:行级(RLAC)与列级(CLAC)
README 将访问控制列为引擎的核心能力之一:RLAC(Row-Level Access Control)与CLAC(Column-Level Access Control)。其实现位于 src/logical_plan/analyze/access_control.rs。
行级访问控制(RLAC)
RLAC 以条件表达式(condition)形式定义在 model 上,条件中可以使用会话属性(以@前缀引用,如@tenant_id)。collect_condition(src/logical_plan/analyze/access_control.rs)会遍历条件表达式并产出两类信息:
- 顶层裸标识符列引用——视为外层 model 的列,被预标记为 required 以免被裁剪;
@name会话属性——无论出现在何处(包括子查询内部)都要收集,供 RLAC 解析时替换。
ConditionVisitor(同一文件第 84-133 行)还承担合法性校验:若条件引用了不属于该 model 的列,会返回The column {} is not in the model {}错误。validate_rlac_rule进一步校验条件语法与 required_properties 是否已在会话属性中定义。
列级访问控制(CLAC)
CLAC 通过validate_clac_rule在注册表阶段生效:如第三节所述,无权访问的列根本不会被注册进WrenDataSource的 schema(src/mdl/mod.rs 中对每个列先做 CLAC 校验再决定是否保留)。因此一个试图访问受限列的查询会得到"列不存在"错误——permission_analyze的职责就是识别这种错误并转换为更友好的权限错误信息(源码注释对此有明确说明)。
六、可运行的完整示例:plan-sql.rs
仓库中的 core/wren-core/wren-example/examples/plan-sql.rs 是官方提供的端到端示例:通过ManifestBuilder编程式构造 MDL(三个 model、两条关系、计算列customer_state),再调用AnalyzedWrenMDL::analyze与transform_sql_with_ctx完成重写。
use wren_core::mdl::builder::{ ColumnBuilder, ManifestBuilder, ModelBuilder, RelationshipBuilder, }; use wren_core::mdl::context::Mode; use wren_core::mdl::manifest::{DataSource, JoinType, Manifest}; use wren_core::mdl::{create_wren_ctx, transform_sql_with_ctx, AnalyzedWrenMDL}; #[tokio::main] async fn main() -> datafusion::common::Result<()> { let manifest = init_manifest(); let analyzed_mdl = Arc::new(AnalyzedWrenMDL::analyze( manifest, Arc::new(HashMap::default()), Mode::Unparse, )?); let sql = "select customer_state from wrenai.public.orders_model"; println!("Original SQL: \n{sql}"); let sql = transform_sql_with_ctx( &create_wren_ctx(None, Some(&DataSource::BigQuery)), analyzed_mdl, &[], HashMap::new().into(), sql, ) .await?; println!("Wren engine generated SQL: \n{sql}"); Ok(()) }这段示例同时展示了几个关键 API 的正确组合:
AnalyzedWrenMDL::analyze(manifest, properties, mode)——在Mode::Unparse下分析 MDL;create_wren_ctx(config, data_source)——创建带正确 UDF 集合的SessionContext(src/mdl/mod.rs),并默认将时区设为 UTC 以避免时间戳推断/比较的时区问题(可被用户配置覆盖);- 传入
Some(&DataSource::BigQuery)时,会话会注册 BigQuery 方言支持的内建 UDF/UDAF/UDWF(通过get_inner_dialect获取,见 src/mdl/context.rs 中create_wren_ctx的实现)。
示例中customer_state是一个跨两跳关系的计算列(orders_model.customers_model.state),引擎需要自动生成orders→customers的 join 才能产出该列——这正是"解析关系链、展开计算列"能力的直观演示。
七、测试与行为验证
单元测试集中在 core/wren-core/core/src/mdl/mod.rs 的mod test中,可作为理解引擎行为的活文档:
test_aliased_model_scan_builds_model_plan_node_once:验证带别名的 model 扫描只会构建一次ModelPlanNode(而不是先建一次被丢弃的通配节点再建一次正式节点),并给出精确的 SQL 快照,例如输入select e.c_custkey from test.test.customer e输出:SELECT e.c_custkey FROM (SELECT customer.c_custkey FROM (SELECT __source.c_custkey AS c_custkey FROM customer AS __source) AS customer) AS etest_access_model:覆盖聚合、join、where、窗口函数等 8 类 SQL 的往返重写,并断言输出 SQL 可执行(assert_sql_valid_executable);test_plan_calculation_without_unnamed_subquery:验证含多跳关系的计算列totalcost被展开为嵌套 join + 聚合的完整 SQL(快照可见RIGHT OUTER JOIN关系链展开);test_uppercase_catalog_schema:验证大小写敏感的 catalog/schema 处理(默认关闭标识符规范化);test_remote_function:从 core/wren-core/core/tests/data/functions.csv 读取远程函数定义并注册,验证add_two、median等 UDF 在重写后被原样保留并生成正确 SQL;test_sync_transform(multi-threadfeature 下):验证同步入口transform_sql可用。
这些测试大量使用insta快照断言,把"重写后的 SQL"固化为可审阅的产物,对引擎行为做了精确到字符串的约束。
八、总结与适用边界
wren-semantic-core的定位可以概括为:一个以 MDL 为输入、以 DataFusion 逻辑计划为中间表示、可输出目标方言 SQL 的语义查询引擎。它同时承担四类职责(与 README 的 "What it does" 一一对应):
- MDL 分析:把 Manifest 解析为 models、columns、metrics(计算列)、relationships 与 views;
- 查询规划:将 SQL 重写为 DataFusion 逻辑计划,解析关系链、展开视图(
ExpandWrenViewRule、ModelAnalyzeRule、ModelGenerationRule); - 访问控制:应用行级(RLAC)与列级(CLAC)规则;
- 优化:类型强制转换(
WrenTypeCoercion)与时间戳简化(TimestampSimplify)等面向 unparse 的裁剪优化。
使用边界与前置条件同样值得注意:transform_sql同步入口依赖multi-threadfeature,WASM 环境请改用异步的transform_sql_with_ctx;analyze_with_url_tables仅支持文件类数据源(LocalFile/Minio/S3/GCS);引擎的 SQL 输出面向目标数据源方言,具体函数能力由WrenDialect与对应数据源的 UDF 集合决定。若要进一步深入,可继续阅读 core/wren-core/README.md 了解整个 wren-core 成员工作区,或查看 core/wren-core/core/CHANGELOG.md 了解版本演进。引擎基于 Apache License 2.0 开源,其语义层设计对构建"可治理的 Text-to-SQL"系统具有直接的参考价值。
【免费下载链接】WrenAIGenBI (Generative BI) for AI agents, an open-source, governed text-to-SQL through an open context layer that turns natural-language questions into trusted dashboards, charts, and SQL across 20+ data sources, such as BigQuery, Snowflake, PostgreSQL, ClickHouse, Amazon Redshift, Databricks and more.项目地址: https://gitcode.com/GitHub_Trending/wr/WrenAI
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考