数据截至 (上游 commit 7830cc746c11)
语义引擎:模型名 SQL → 受治理物理 SQL
30 秒导读: 这是 Wren 整个项目名所指的核心价值。你(或 agent)写一句看起来很普通的 SQL——
SELECT customer_name, total FROM orders——但orders不是真实数据库里的表,而是 MDL(语义建模语言)里定义的一个模型。语义引擎负责把这句"模型名 SQL"翻成真实库表上、 带 join、带权限过滤的物理 SQL,再交给对应数据源执行。改写的每一步都是治理挂钩点: 越权的表挡掉、禁用的函数挡掉、看不到的行和列在 SQL 里就被裁掉。
本章是这套参考里最深的一章。它由一个薄薄的 Python 门面(WrenEngine)一路追到 Rust
语义引擎(wren-core,基于 Apache DataFusion 的 Canner fork)。记忆检索见
04-memory-and-retrieval.md,仪表盘生成见
05-genbi-dashboard-deploy.md;MDL 模型本身的结构见
02-context-layer-mdl.md。
1. 这是什么(零基础也能懂)
一句话定义: 语义引擎是一个"SQL 编译器",输入是针对语义模型写的 SQL,输出是针对 真实数据库表、且已经施加了访问控制的 SQL。
它解决什么问题
想象你有一张真实的 public.orders_2024 表,列名是 cust_id、amt_cents,还要 join 一张
customers 表才能拿到客户名。让一个 AI agent(甚至人)每次都记住这些物理细节,既容易写错、
又没法统一加权限。
MDL 的做法是:定义一个叫 orders 的模型,里面声明"我的 customer_name 其实来自
join customers 后的 name 列"、"我的 total 是 amt_cents / 100"。之后所有人只写:
-- 用户/agent 写的:干净、稳定、跟物理表解耦
SELECT customer_name, total FROM orders WHERE total > 100
语义引擎把它编译成真正能在 Postgres 上跑的东西——join 补上、列名换回物理名、还按当前用户的 身份自动加上"只能看自己区域的订单"这种过滤。
它给谁用
- AI agent / MCP 客户端——写业务语义的 SQL,不必了解物理 schema,天然被治理兜住。
- 数据平台——一处定义模型 + 权限,22 种数据源上行为一致。
用起来什么样
WrenEngine 是 Python 侧唯一入口(core/wren/src/wren/engine.py:49):
# 示意,非源码(改编自 engine.py 顶部 docstring)
engine = WrenEngine(
manifest_str="<base64 编码的 MDL JSON>", # 语义模型
data_source=DataSource.postgres, # 目标数据源
connection_info={"host": "localhost", ...},
)
# 只做改写,不碰数据库:输出目标方言的物理 SQL
planned_sql = engine.dry_plan("SELECT * FROM orders")
# 改写 + 执行,拿回 Arrow 表
table = engine.query("SELECT * FROM orders", limit=100)
一句话直觉
把 MDL 模型当"视图",但这个视图不是数据库里的对象,而是一份可版本化、可施加权限的规格; 语义引擎就是在查询时把视图现场展开、并顺手把权限焊进 SQL 的编译器。
本节到此不涉及任何底层代码——只要知道"它把模型名 SQL 编译成受治理的物理 SQL"就够了。
2. 顶层全景(它大概怎么转)
一次 dry_plan/query 从 Python 门面下钻到 Rust 引擎、再落到连接器,主线如下。
怎么读这张图: 从上往下是一次调用的时间顺序;左边是 Python(core/wren),中间过
PyO3 桥进入 Rust(wren-core),右边是最终执行。带 ★ 的方框是治理挂钩点。
┌─────────────────────────────────────────────┐
用户/agent SQL ─▶│ WrenEngine._plan (engine.py) │ Python 门面
(目标方言) │ │
│ ① sqlglot 按目标方言解析 │
│ ② ★ 策略校验 validate_sql_policy │ policy.py
│ ③ 大小写归一 resolve_model_name │
│ ④ 抽取最小 manifest extract_by │──┐
└─────────────────────────────────────────────┘ │ PyO3
┌─────────────────────────────────────────────┐ │
│ CTERewriter.rewrite (cte_rewriter.py) │◀─┘
│ 逐个模型 → session.transform_sql │──┐
│ 展开结果作为 CTE 注入,视图逐字注入 │ │ PyO3
└─────────────────────────────────────────────┘ │
┌─────────────────────────────────────────────┐ │
│ wren-core (Rust) transform_sql_with_ctx │◀─┘ 语义引擎
│ DataFusion 三条 AnalyzerRule 依次跑: │
│ ExpandWrenViewRule 展开视图 │
│ ModelAnalyzeRule TableScan→ModelPlan │
│ ★ + RLAC/CLAC 行列级裁剪 │ access_control.rs
│ ModelGenerationRule 展开成底层表 join │
│ Unparser → 目标方言物理 SQL │
└ ─────────────────────────────────────────────┘
┌─────────────────────────────────────────────┐
Arrow 表 ◀───────│ get_connector → PostgresConnector.query │ 连接器
│ (22 种数据源;原生驱动执行) │ connector/
└─────────────────────────────────────────────┘
各部件一句话职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
WrenEngine | 门面:编排解析→校验→抽取→改写→执行 | core/wren/src/wren/engine.py |
validate_sql_policy | 严格模式 + 函数黑名单的治理原语 | core/wren/src/wren/policy.py |
ManifestExtractor.extract_by | 只抽出 SQL 用到的模型,缩小 manifest | core/wren-core-py/src/extractor.rs |
CTERewriter | 逐模型改写、把展开结果作为 CTE 注入 | core/wren/src/wren/mdl/cte_rewriter.py |
PySessionContext.transform_sql | PyO3 桥:把单模型 SQL 交给 Rust 展开 | core/wren-core-py/src/context.rs |
transform_sql_with_ctx | Rust 主流程:建计划→跑规则→反解析成 SQL | core/wren-core/core/src/mdl/mod.rs |
三条 AnalyzerRule | 视图展开 / 模型分析 / 生成 join,含 RLAC/CLAC | core/wren-core/core/src/logical_plan/analyze/ |
get_connector | 按数据源分派到具体连接器 | core/wren/src/wren/connector/factory.py |
主线走一遍(高层): 一句模型名 SQL 进来 → Python 用 sqlglot 解析并做策略校验 → 抽出最小
manifest → 对每个被引用的模型,单独发一句 SELECT ... FROM 模型 给 Rust 展开成物理 SQL →
把展开结果当 CTE 塞回原查询 → 原查询里对模型的引用于是绑定到这些 CTE → 输出目标方言的物理
SQL → 连接器执行返回 Arrow。
3. 核心原理(逐个机制,由浅入深)
3.1 门面 _plan:五步链路
_plan(engine.py:177,由 dry_plan/query/dry_run 共用)是整章的骨架。它做五件事,
每件都对应一个下钻方向。
它要解决的小问题: 把一句"目标方言、含模型名"的 SQL,安全地喂给只认单模型的 Rust 引擎。
思路: 不把整句 SQL 直 接丢给 Rust,而是先在 Python 侧用 sqlglot 拆解、校验、缩小范围, 再逐模型下钻。这样治理(策略校验)发生在最外层,且能把昂贵的 Rust 分析限制在最小 manifest 上。
五步依次是(engine.py:186-251):
| 步 | 做什么 | 关键调用 |
|---|---|---|
| ① 解析 | 按目标方言(Postgres/BigQuery…)解析成 AST | parse_one(sql, dialect=...) engine.py:189 |
| ② 校验 | 总是调用:只读语句检查不依赖开关,严格模式 / 函数黑名单再挡非法表与函数 | validate_sql_policy(...) engine.py:203 |
| ③ 归一 | 把表引用解析成 manifest 里的规范模型名 | resolve_model_name(...) engine.py:216 |
| ④ 抽取 | 只抽出用到的模型,得到最小 manifest | extractor.extract_by(tables) engine.py:220 |
| ⑤ 改写 | 建 SessionContext + CTERewriter.rewrite | rewriter.rewrite(sql) engine.py:247 |
⑤ 之后还有一道对改写产出的复核 validate_planned_sql(engine.py:250,实现在
policy.py:259):改写会把 MDL 视图定义和模型 ref_sql inline 进计划结果,输入侧的检查
看不见这些后加进来的语句,所以产出让连接器执行前要再查一次只读(解析失败则 fail-open 放行)。
一个关键设计:queryable_names = 模型名 ∪ 视图名(engine.py:197)。视图也是 MDL 定义的
对象,所以严格模式放行视图引用,而 extract_by 会把视图连同它 join 的模型一起抽进来。
还有一个务实的容错分支(engine.py:224-233):抽取最小 manifest 若失败,只要没开严格模式、
也没配函数黑名单,就回退用完整 manifest——治理没被绕过(那两种模式下失败会抛
INVALID_SQL,且 ② 的只读校验此时已经跑过),只是优化(缩小 manifest)可以放弃。
真实实现,engine.py:219-233:
# 真实源码节选(engine.py:_plan)
extractor = get_manifest_extractor(self.manifest_str)
manifest = extractor.extract_by(tables)
effective_manifest = to_json_base64(manifest)
# ...except Exception:
if self._config.strict_mode or self._config.denied_functions:
raise WrenError(ErrorCode.INVALID_SQL, ..., phase=ErrorPhase.SQL_PLANNING, ...)
effective_manifest = self.manifest_str # 回退全量 manifest
这段把"抽取失败"和"治理必须成立"两件事分开:治理开着时,任何异常都变成结构化错误往上抛; 治理关着时,退回全量 manifest 让改写继续。
3.2 治理原语:validate_sql_policy
对应 README 里的 "governed execution primitives"(README.md:184)。这是在 SQL 真正被
展开、执行之前的第一道闸。policy.py:280 的 validate_sql_policy 只有两个开关:
strict_mode→_check_data_readers+_check_tables:先全树挡掉从 manifest 外读数据的 表值函数,再要求表引用只指向 manifest 内的对象;denied_functions→_check_functions:黑名单里的函数一律挡掉。
严格模式 _check_tables(policy.py:332)
它遍历 AST 里所有 exp.Table,每个都要么解析成一个模型/视图名,要么是当前作用域可见的
CTE,否则抛 MODEL_NOT_FOUND。可见 CTE 的判定靠 _visible_cte_names(policy.py:311)
——沿 AST 往上走,收集每一层 WITH 定义的名字。
严格模式对危险函数按上游 issue #2409 的分类区别对待:
- 数据/文件读取器(
read_csv/dblink/postgres_scan/ …):由_check_data_readers(policy.py:498)在 AST 的任何位置挡掉——不只FROM/JOIN源位置,投影、WHERE 子查询、嵌套参数(如UNNEST(read_csv(...)))里的也逃不掉,堵路径穿越 / SSRF / 数据外泄; - 合成生成器(
generate_series/sequence/range):不读 manifest 外数据,但无界范围是 DoS 向量,所以源位置默认挡、运维可在 config 里用allowed_source_functions按名单放行 (_is_allowed_generator,policy.py:477)。
表值函数(TVF)没有 exp.Table 节点。read_csv(...)、generate_series(...) 这类会直接
作为 FROM/JOIN 的源出现,所以额外扫 exp.From, exp.Join(policy.py:379),源是
exp.Func 就挡(生成器按上述白名单豁免)。带别名的 TVF(generate_series(1,10) AS t(x))
解析成 this 为内层函数的 exp.Table,判定同样下探到内层函数(policy.py:340-358)。
但行展开算子(UNNEST/FLATTEN/EXPLODE,policy.py:38 的 _ROW_EXPANSION_FUNCS)
是例外:它们重构的是已经在查询作用域内的数组/结构列(比如某个受治理模型的 orders.items),
并不从 manifest 外读数据,所以放行。
真实实现,policy.py:340-358:
# 真实源码节选(policy.py:_check_tables)
if not name:
# 没有名字的 Table 节点就是表值函数(read_csv() 等)。
# 数据读取器已被 _check_data_readers 全树挡掉;
# 生成器按 allowed_source_functions 白名单放行。
inner = table.this
if isinstance(inner, exp.Func) and _is_allowed_generator(
inner, allowed_source_functions
):
continue
sql_text = table.sql()
if sql_text:
raise WrenError(ErrorCode.MODEL_NOT_FOUND,
f"Table-valued function '{sql_text}' is not allowed. ...",
phase=ErrorPhase.SQL_POLICY_CHECK)
函数黑名单 _check_functions 与 _canonical_names(policy.py:413)
难点在于:sqlglot 会把同一个函数名映射到不同的 AST 子类,取决于方言。比如 version()
在 postgres/mysql/duckdb/trino/clickhouse 下被规约成 exp.CurrentVersion,而在
tsql/oracle/bigquery/snowflake 下仍是 exp.Anonymous。用户黑名单里只写 "version" 的话,
只会命中匿名那种。
_canonical_names(policy.py:413,带 @lru_cache;旧名 _canonical_denied 现在只是
它的别名,policy.py:455)的解法:把每个名字在十种方言下、用三种实参形态
(name() / name('x') / name(1, 2))各解析一遍——有的函数带参数才会落到具体子类
(如 duckdb 的 read_csv('x') → exp.ReadCSV)——把它落到的每个具体类的 key 都收进规范
集合(policy.py:429-452)。这套规范化同时服务黑名单、数据读取器黑名单和生成器白名单,
无论用户 SQL 被哪个方言解析、sqlglot 是否把它重分类,都能命中。
3.3 大小写归一:resolve_model_name
它要解决的小问题: SQL 的标识符大小写规则很微妙——带引号严格区分大小写,不带引号各方言
折叠规则不同。而 Rust 侧的 extract_by 是大小写敏感的,必须先把用户写的 Orders / "orders"
归一到 manifest 里真实的模型名。
resolve_model_name(policy.py:152)是贯穿改写器、策略检查、抽取器三处的同一条 SQL 约定:
# 真实源码(policy.py:resolve_model_name,已精简注释)
if name in model_set: # 精确匹配优先(引号/非引号都先试)
return name
if quoted: # 带引号:严格,不做大小写回退
return None
name_lower = name.lower() # 非引号:回退到大小写不敏感扫描
for candidate in model_set:
if candidate.lower() == name_lower:
return candidate
return None
engine.py:210-217 用它把 AST 里每个表名归一后,才交给 extract_by。
3.4 CTE 改写器:把模型变成 CTE
这是 Python 侧最重的一块(cte_rewriter.py,1124 行)。核心思路一句话:
从不改用户的 SQL,只在前面追加 CTE。 对每个被引用的模型,单独发一句
SELECT col1, col2 FROM "模型" 给 Rust 展开成物理 SQL,再把展开结果作为一个与用户所写
同名的 CTE 注入;于是用户原句里对模型的引用,自然绑定到这个 CTE。
rewrite(cte_rewriter.py:236)的主流程:
用户 SQL ──parse──▶ AST
│
├─ 收集用户自定义 CTE 名(避免误当模型) _collect_user_cte_names
├─ 用 qualify_columns 解析出每个模型用到哪些列 _collect_model_columns
├─ 收集被引用的视图 _collect_view_refs
│
├─ 逐模型:SELECT <列> FROM "模型"
│ └─▶ session.transform_sql ──▶ 物理 SQL ──▶ 包成 CTE _build_model_ctes
├─ 逐视图:把 view.statement 原样包成 CTE(不经 Rust) _build_view_ctes
│
└─ 把 model CTE + view CTE 追加到 WITH 前部 _inject_ctes
模型 vs 视图的关键区别(cte_rewriter.py 顶部 docstring):模型由 wren-core 展开;
而视图的 statement 本身就是引用模型的原生 SQL,所以逐字注入成 CTE(永不发给 Rust),
它引用的模型则作为 model CTE 排在它前面。这样视图保持它被编写时的可执行 SQL 形态。
_build_model_ctes(cte_rewriter.py:844)里对每个模型只发最小的一句:
| 用户怎么引用模型 | 发给 Rust 的 SQL | 为什么 |
|---|---|---|
SELECT * | SELECT * FROM "model" | 让 wren-core 控制列可见性(CLAC) |
| 引用了具体列 | SELECT "model"."c1", ... FROM "model" | 只展开需要的列 |
只用到行(如 COUNT(*)) | SELECT 1 FROM "model" | 只需要行,不需要列 |
方言映射由 get_sqlglot_dialect(cte_rewriter.py:44)负责,比如 canner→trino、
mssql→tsql、文件源→duckdb。
大量代码在处理一个真实世界的脏问题:大小写折叠。不同方言对未加引号标识符的折叠规则不同
(Postgres 折小写、Oracle/Snowflake 折大写、BigQuery/DuckDB 连反引号也折小写),而模型可能声明
了仅大小写不同的列(Year 和 year)。只有 _CASE_SENSITIVE_COLUMN_DIALECTS
(cte_rewriter.py:62,即 postgres/oracle/snowflake/clickhouse)物理上能区分这种列;其它方言
上这种模型直接在构建期被 INVALID_MDL 拒掉(_raise_case_collision,cte_rewriter.py:211)。
这套"大小写敏感路径"仅在 manifest 真的 含大小写冲突列时才启用,让绝大多数现有 manifest 走原来
成熟的大小写不敏感路径。
3.5 PyO3 桥:Python 如何调 Rust
mdl/__init__.py 是薄薄一层,把 wren_core(PyO3 编出来的扩展模块)的能力包装出来:
| Python 包装 | 底层 Rust | 作用 |
|---|---|---|
get_session_context(mdl/__init__.py:9,装饰器 @lru_cache(maxsize=32)) | wren_core.SessionContext | 建/复用会话上下文 |
get_manifest_extractor(:26) | wren_core.ManifestExtractor | 抽最小 manifest |
to_json_base64(:30) | wren_core.to_json_base64 | manifest 对象 → base64 |
transform_sql(:34) | session.transform_sql | 单句 SQL 展开 |
get_session_context 上的 @lru_cache(maxsize=32) 是个重要细节:相同的
(manifest_str, function_path, properties, data_source) 元组会复用同一个 SessionContext
——所以约定不要去 mutate 会话状态(见 core/wren/.claude 备注)。缓存从进程级
@cache 改成有界 LRU 是刻意的:缓存键含 engine.py extract_by 抽出的每查询最小
manifest,无界缓存会随不同表子集无限增长(见 mdl/__init__.py:15-18 docstring)。
Rust 侧的 PySessionContext(wren-core-py/src/context.rs:84)在构造时(context.rs:121)就一次性
analyze 好 MDL、并预建了两种上下文(Unparse 反解析用、LocalRuntime 本地执行用)。
transform_sql(wren-core-py/src/context.rs:256)把 SQL 转交给 mdl::transform_sql_with_ctx,
期间释放 GIL、跑在进程级共享的 Tokio runtime 上(fork 后惰性重建)——同一个 context
可以并发调用,每次调用套一份私有 catalog 快照(context.rs:75-83 docstring)。
4. 深入实现:Rust 语义引擎
进入 wren-core 后,SQL 的展开完全靠 DataFusion 的 AnalyzerRule 机制。这是本章最底层的一段。
4.1 主流程 transform_sql_with_ctx
mdl/mod.rs:482。它的骨架:
register 远程函数
→ apply_wren_on_ctx(ctx, mdl, Mode::Unparse) // 挂上 Wren 的分析规则
→ ctx.state().create_logical_plan(sql) // DataFusion 建逻辑计划
→ ctx.state().optimize(&plan) // 跑分析/优化规则(展开在这里发生)
→ Unparser::plan_to_sql(&analyzed) // 把计划反解析回 SQL
→ 去掉 catalog.schema 前缀,得到物理 SQL
关键在于:模型的展开不是字符串替换,而是发生在逻辑计划层——先把 SQL 建成一棵计划树, 让分析规则在树上做变换,再反解析回目标方言的 SQL。这也是为什么 Wren 能跨 22 种方言输出正确 SQL。
一个用户友好设计:若建计划失败,会走 permission_analyze(mdl/mod.rs:553)用
Mode::PermissionAnalyze 再分析一遍,专门判断"这个报错其实是权限拒绝"——因为不可见的列压根
不会被注册,正常流程里只会报"列不存在",而这个二次分析能把它翻译成更友好的权限错误。
4.2 三条分析规则(展开的核心)
规则清单在 mdl/context.rs,按 Mode 不同而不同。反解析用的
analyze_rule_for_unparsing(mdl/context.rs:255)依次挂三条 Wren 规则(顺序要紧):
怎么读: 每条规则接收上一条产出的计划树, 做一次变换往下传。
DataFusion 逻辑计划(含对模型的 TableScan)
│
① ExpandWrenViewRule 把视图的 TableScan 换成它的子计划
(expand_view.rs:12) —— 必须第一个跑
│
② ModelAnalyzeRule 自底向上、深度优先地收集每个模型
(model_anlayze.rs:48) 需要哪些列,把 TableScan 变成 ModelPlanNode;
│ 并在此施加 RLAC/CLAC(下节)
│
③ ModelGenerationRule 把 ModelPlanNode 真正展开成
(model_generation.rs:29) "底层表 + relationship join" 的计划
ModelAnalyzeRule 的 doc(model_anlayze.rs:35-47)写明三步:①分析作用域、收集模型和访问过的
表所需的列(自底向上、深度优先);②按作用域分析生成 ModelPlanNode;③去掉 Wren 的
catalog/schema 前缀并刷新 schema(自顶向下)。作用域追踪由 scope.rs 的 Scope/ScopeManager
承担(scope.rs:26)——它记录每个查询作用域里"数据集需要哪些列",父作用域的关系能被子作用域
访问,子作用域也能把所需列上报给父作用域。
ModelGenerationRule 真正把模型摊平成 join。join 结构由 RelationChain(relation_chain.rs:35)
表达——它是一串"模型 + join 类型 + 条件"的链,物理布局形如 (((Model3, Model2), Model1), Nil)。
plan(relation_chain.rs:139)递归地把这条链构造成嵌套的 DataFusion join 计划。
4.3 行/列级访问控制(RLAC/CLAC)
access_control.rs 是治理焊进 SQL 的地方——README 在 "Governed execution, reviewable
context" 里点名的 "row/column-level security and access control" 在上游定位为 Cloud /
self-hosted 商业版能力(README.md:56、README.md:225),但机制本身在 OSS 引擎里就是这几条分析规则。
行级(RLAC): 模型可以带一条"行级访问控制"规则,它是一段条件表达式(可含
@session_id 这种会话属性)。collect_condition(access_control.rs:52)解析这个条件,提取
①条件顶层引用的裸列(预标为"必需",避免被裁剪)②条件里(含子查询内)引用的会话属性——后者
要在 RLAC 解析期被替换。这条件最终变成 SQL 里的一个 WHERE,把当前身份看不到的行挡在计划里。
一个安全细节:ModelAnalyzeRule 做 RLAC 解析时用一个环检测栈检测相互引用的环(A 的 RLAC
引用 B,而 B 的 RLAC 又引用 A)——ModelStack(model_anlayze.rs:59,RefCell<HashSet>,
按每次 analyze 调用分配、随递归下传),配 RAII 的 ModelStackGuard 保证出栈。
(它不再是挂在规则字段上的 Arc<Mutex<...>>:同一 context 共享一条 Send + Sync 规则,
RefCell 存在字段上会编译不过。)
列级(CLAC): validate_clac_rule(access_control.rs:534)在生成模型计划时判断某列对当前
会话是否可见。不可见的列根本不会被注册进计划——所以 SELECT * 时被裁掉、显式选它时会因
"列不存在"而失败(再由 §4.1 的 permission_analyze 翻译成权限错误)。这就是为什么 §3.4 里
SELECT * 要原样发给 Rust:把列可见性的决定权留给 wren-core。
4.4 抽取最小 manifest:extract_by
extractor.rs:45。给定 SQL 用到的数据集名字列表,extract_manifest(extractor.rs:117)剔掉
没用到的模型/视图,但保留关联:若一个模型通过 relationship 关联到另一个数据集,两个都留、
它们之间的 relationship 也留(extractor.rs:42-47 的 doc)。视图会连带把它引用的模型抽进来
(extract_views)。这一步把后续昂贵的 Rust 分析限制在最小范围内。
resolve_used_table_names(extractor.rs:37/50)则演示了它如何解析表引用:关掉标识符归一
(enable_ident_normalization = false),再用 DataFusion 的 resolve_table_references,并过滤掉
catalog/schema 与本 manifest 不匹配的表。
PyO3 绑定的暴露面在 lib.rs:13-29:PySessionContext、PyManifestExtractor、Manifest/Model/
RowLevelAccessControl 等类型,以及 to_json_base64、validate_rlac_rule、cube_query_to_sql
等函数,全在这个 #[pymodule] 里注册给 Python。
5. 执行落地:连接器与配置
改写完成后,query/dry_run 才需要真数据库。这一层的入口是 get_connector。
5.1 连接器工厂
connector/factory.py:46。它按 DataSource 查一张注册表(_REGISTRY,factory.py:6)找到连接器
模块,动态 import,再调其 create_connector。几处映射体现了"多对一"的复用:
| 数据源 | 实际连接器模块 | 说明 |
|---|---|---|
doris | wren.connector.mysql | MySQL 兼容协议 |
canner | wren.connector.postgres | Postgres 线协议 |
local_file/s3_file/minio_file/gcs_file | wren.connector.duckdb | 文件源统一走 DuckDB |
导入失败会翻译成友好的 NOT_IMPLEMENTED 错误并提示 pip install wren[<extra>]
(factory.py:56-62)。所有连接器都实现同一个抽象基类 ConnectorABC(connector/base.py:61):
query(sql, limit)、dry_run(sql)、close() 三个方法。基类文件还沉淀了两个所有连接器共用的
守门工具:strip_trailing_semicolon(base.py:11,剥掉结尾分号,避免子查询包裹/EXPLAIN 形态
被引擎拒)和 coerce_limit(base.py:23,把用户传的 limit 统一校验成非负 int,非法值在拼
LIMIT 前就报错)。
一处与仓库内 CLAUDE.md 的出入(诚实记录): 仓内
core/wren/.claude/CLAUDE.md称连接器 "wrap 一个 Ibis backend"。但实际读源码,connector/下只有canner.py还 import ibis;postgres.py(psycopg)、snowflake.py(snowflake-connector)、oracle.py(oracledb)、trino.py、clickhouse.py(clickhouse-connect)、athena.py(pyathena)等都已改成原生驱动, docstring 明写 "bypasses the ibis backend"。以源码为准。
以 Postgres 为例(connector/postgres.py:276):PostgresConnector 直接用 psycopg3 执行,并用一张
手写的 PG OID → Arrow 类型表(_PG_OID_TO_ARROW,postgres.py:28)把游标结果转成 PyArrow。
query(postgres.py:304)先 coerce_limit 校验、strip_trailing_semicolon 剥分号,limit
再包一层 SELECT * FROM (<sql>) AS _sub LIMIT n(postgres.py:310);dry_run 则是
LIMIT 0(postgres.py:329)。执行异常统一裹成带 dialectSql 元数据的 WrenError。
5.2 结构化错误:WrenError
model/error.py。治理和执行途中的每一步失败都变成一个 WrenError(error.py:49),带三样东西:
ErrorCode(error.py:9):INVALID_SQL、MODEL_NOT_FOUND、BLOCKED_FUNCTION、DATABASE_TIMEOUT等枚举,给机器判别;ErrorPhase(error.py:34):失败发生在哪一阶段——SQL_POLICY_CHECK、SQL_PLANNING、SQL_EXECUTION、MDL_EXTRACTION等,让 agent 知道"是策略挡的还是库报的";metadata:常带dialectSql(error.py:5),把出错时的 SQL 附上,方便定位。
这套结构化错误正是 README 说的 "structured errors with hints"(README.md:55)——agent 拿到的不是
一坨字符串,而是能据以决定"改 SQL 还是换库"的分类信号。
5.3 连接配置与密钥:profile.py
连接信息不写死在代码里,而是走命名 profile(~/.wren/profiles.yml)。profile.py 管这些
profile 的增删查改,以及一个关键的安全能力:密钥延迟展开。
profile 值里可以写 ${PG_PASSWORD} 这种占位符;它们只在连接时由
expand_profile_secrets(profile.py:122)解析,存盘的 YAML 始终保留占位符,所以
wren profile debug 永远不会打印出真实密钥。展开时:
_SecretTemplate(profile.py:30)刻意只匹配${UPPER_SNAKE_CASE},让真实密码里的 小写${foo}序列原样保留;_ensure_env_loaded(profile.py:45)按$CWD/.env→ 项目根.env→~/.wren/.env的顺序 合并环境变量,且已导出的环境变量永远优先(override=False);- 引用了未设置的变量会抛
MissingSecretError(profile.py:38),给出可操作的提示。
存盘走原子写 + 0600 权限(_save_raw,profile.py:181),debug_profile(profile.py:375)还会
打码敏感字段——以连接字段注册表里标 SecretStr 的字段为准(_registry_sensitive_keys,
profile.py:325;兜底再加 password/credentials/secret/token 等名称启发式),并递归走嵌套
dict/list(_mask_obj,profile.py:358),因为 kwargs/settings 里也可能藏密钥。