PostgreSQL QuerySpan 迁移到 Omni:Bytebase 以语义分析替代 ANTLR 的列血缘提取重构实践
2026/9/14 18:09:08 网站建设 项目流程

PostgreSQL QuerySpan 迁移到 Omni:Bytebase 以语义分析替代 ANTLR 的列血缘提取重构实践

【免费下载链接】bytebaseDatabase governance built for humans and agents — controlling changes and access across every major database.项目地址: https://gitcode.com/GitHub_Trending/by/bytebase

导读

本文基于 docs/plans/2026-03-24-pg-query-span-omni-migration.md 这一实施计划,完整拆解 Bytebase 将 PostgreSQL QuerySpan(查询跨度:访问表集合 + 结果列血缘)提取从基于 ANTLR 的语法树遍历,迁移到 omni 语义分析基础设施的全过程。本次改造以约 670 行新增代码替换约 3,964 行遗留代码,净减约 3,300 行,同时保证既有 46 个 query_span 与 31 个 query_type YAML 用例全部通过。读完本文,你将掌握 omni 的AnalyzeSelectStmt语义分析管线、基于Query结构的列血缘 walker 设计,以及如何用"逐步替换 + 全量 YAML 回归"的策略安全完成一次高危核心模块迁移。


一、背景:为什么要把 ANTLR 换成 omni

当前状态(迁移前)

Bytebase 的 PostgreSQL 解析器在迁移前依赖 ANTLR 语法树做手工遍历,相关代码分散在四个文件中:

文件职责体量
query_span.go入口,调用 ANTLR 提取器
query_span_extractor.goANTLR 树遍历提取(约 70 个方法)3,868 行
access_tables_antlr.goANTLR listener 提取被访问的表96 行
query_type.go查询类型分类已迁移到 omni AST,无需改动
access_tables.go基于 omni 的ExtractAccessTables()已迁移,无需改动

ANTLR 方案的核心问题:SQL 是声明式语言,语法树只描述"写成了什么样",不描述"实际访问了什么"。SELECT a FROM t1 JOIN t2SELECT * FROM v(v 是视图)在语法上结构相同,但列血缘完全不同。要用语法树还原语义,只能在 extractor 里手工模拟解析器的名称解析、作用域、星号展开等逻辑——这正是 3,868 行代码的由来,且每支持一种新语法(CTE、集合运算、视图、窗口函数)都要打补丁。

目标状态(迁移后)

文件职责体量
query_span.go入口,调用 omni 提取器
query_span_omni.go新建:约 400–600 行,使用AnalyzeSelectStmt()+ 血缘 walker新增
access_tables_antlr.go删除(由access_tables.go取代)-96
query_span_extractor.go迁移完成后删除-3,868

从当前仓库源码结构看,本次迁移已经落地:backend/plugin/parser/pg目录下已不存在query_span_extractor.goaccess_tables_antlr.go,入口文件 query_span.go 已完全基于 omni 实现(最终文件约 1,800 行,除血缘 walker 外还包含函数体分析、fallback 列提取等健壮性逻辑)。

范围排除

PL/pgSQL 函数体分析单独跟踪在 BYT-9082 工单中。迁移期间,函数调用将回退到既有 ANTLR 函数分析逻辑,直到 BYT-9082(omni 内置 PL/pgSQL 解析器)完成。


二、核心架构:omni 的语义分析管线

技术栈

  • 语言:Go
  • 核心依赖github.com/bytebase/omni(见 go.mod,当前锁定版本v0.0.0-20260912023254-4574e69bb9f1
  • 元数据来源:Bytebase 现有的数据库元数据基础设施(backend/store/modelbackend/plugin/parser/base

三步式分析流程

  1. 解析:omni 的ParsePg()把 SQL 解析成 AST(入口封装见 omni.go 的func ParsePg(sql string) ([]omnipg.Statement, error))。
  2. 语义分析:调用catalog.Catalog.AnalyzeSelectStmt(selStmt),产出带解析信息的Query结构——列引用(VarExpr)已解析到具体的RangeTableEntry,附带类型信息与来源追踪(provenance)。
  3. 血缘提取:在分析后的Query树上行走,把每个TargetEntry(结果列)映射到一组ColumnResource{Database, Schema, Table, Column}

Schema 元数据如何进入 omni catalog

计划中的原始方案复用了walk_through_omni.go中已经验证的模式:

schema.GetDatabaseDefinition(Engine_POSTGRES, ctx, metadataProto) → schemaDDL catalog.New() + SetSearchPath() + catalog.Exec(schemaDDL, ContinueOnError)

即先把 Bytebase 的元数据 proto 序列化成 DDL,再"回放"进 omni 的 catalog。

实现演进:当前仓库的 query_span.go 中的initCatalog()已不再走"生成 DDL 再回放"的路径,而是直接调用e.cat.LoadMetadata(ctx, meta.GetProto(), catalog.LoadMetadataOptions{})将元数据 proto 直接载入 catalog,并会对report.Degraded/report.Missing(降级为 stand-in 或缺失的对象)打slog.Debug日志。生成 DDL 的能力本身仍保留在 get_database_definition.go(由 schema.go 统一暴露GetDatabaseDefinition),供其他场景使用。

关键数据结构(omni/pg/catalog 包)

  • Query:分析后的查询,含TargetList(结果列)、RangeTable(范围表)、CTEListSetOp/LArg/RArg(集合运算)、JoinTree等。
  • TargetEntry:一个结果列,含ExprResNameResJunk(系统辅助列标记)。
  • VarExpr:列引用表达式,通过RangeIdx(指向RangeTable的索引)+AttNum(1 起始的列序号)唯一定位。
  • RangeTableEntry:FROM 中的每个来源,Kind区分物理表(RTERelation)、子查询(RTESubquery)、CTE(RTECTE)、函数(RTEFunction)、JOIN(RTEJoin)。

三、任务拆解:11 步完成迁移

计划将整个迁移拆成 11 个可独立提交的任务,每个任务都以"实现 → 测试 → 提交"闭环推进。下面按阶段分组展开。

阶段一:脚手架与入口接线(Task 1)

新建query_span_omni.go,定义提取器结构与构造函数:

package pg import ( "context" "github.com/pkg/errors" "github.com/bytebase/omni/pg/ast" "github.com/bytebase/omni/pg/catalog" storepb "github.com/bytebase/bytebase/backend/generated-go/store" "github.com/bytebase/bytebase/backend/plugin/parser/base" "github.com/bytebase/bytebase/backend/plugin/schema" "github.com/bytebase/bytebase/backend/store/model" ) // omniQuerySpanExtractor extracts query span using omni's semantic analysis. type omniQuerySpanExtractor struct { ctx context.Context gCtx base.GetQuerySpanContext defaultDatabase string searchPath []string metaCache map[string]*model.DatabaseMetadata cat *catalog.Catalog } func newOmniQuerySpanExtractor( defaultDatabase string, searchPath []string, gCtx base.GetQuerySpanContext, ) *omniQuerySpanExtractor { if len(searchPath) == 0 { searchPath = []string{"public"} } return &omniQuerySpanExtractor{ defaultDatabase: defaultDatabase, searchPath: searchPath, gCtx: gCtx, metaCache: make(map[string]*model.DatabaseMetadata), } }

要点说明:

  • searchPath为空时默认["public"],与 PostgreSQL 的默认搜索路径一致。
  • metaCache做数据库元数据的惰性缓存,避免同一次分析重复拉取。
  • 入口GetQuerySpan重接线:先通过gCtx.GetDatabaseMetadataFunc取元数据,读取meta.GetSearchPath()作为搜索路径(schema参数非空时以指定 schema 覆盖),再构造 omni 提取器调用getQuerySpan(ctx, stmt.Text)。当前 query_span.go 中该入口对 PostgreSQL 与 CockroachDB 两个引擎统一注册(base.RegisterGetQuerySpan(storepb.Engine_POSTGRES, GetQuerySpan)storepb.Engine_COCKROACHDB)。

提交:

git add bytebase/backend/plugin/parser/pg/query_span_omni.go bytebase/backend/plugin/parser/pg/query_span.go git commit -m "feat(pg): scaffold omni-based QuerySpan extractor"

阶段二:catalog 元数据加载(Task 2)

实现 catalog 初始化。计划版本使用GetDatabaseDefinition → Exec(schemaDDL)

func (e *omniQuerySpanExtractor) getDatabaseMetadata(database string) (*model.DatabaseMetadata, error) { if meta, ok := e.metaCache[database]; ok { return meta, nil } _, meta, err := e.gCtx.GetDatabaseMetadataFunc(e.ctx, e.gCtx.InstanceID, database) if err != nil { return nil, errors.Wrapf(err, "failed to get database metadata for database: %s", database) } e.metaCache[database] = meta return meta, nil } // initCatalog creates an omni catalog loaded with the database schema. func (e *omniQuerySpanExtractor) initCatalog() error { meta, err := e.getDatabaseMetadata(e.defaultDatabase) if err != nil { return err } schemaDDL, err := schema.GetDatabaseDefinition( storepb.Engine_POSTGRES, schema.GetDefinitionContext{}, meta.GetProto(), ) if err != nil { return errors.Wrap(err, "failed to generate schema DDL") } e.cat = catalog.New() e.cat.SetSearchPath(e.searchPath) if schemaDDL != "" { if _, err := e.cat.Exec(schemaDDL, &catalog.ExecOptions{ContinueOnError: true}); err != nil { return errors.Wrap(err, "failed to load schema into catalog") } } return nil }

验证方式:构造一份 metadata proto,调用initCatalog(),然后通过cat.GetRelation("public", "t")确认表能查到。测试命令:

go test -v -count=1 -run ^TestOmniCatalogLoading$ github.com/bytebase/bytebase/backend/plugin/parser/pg

从当前源码看,initCatalog的最终形态改用了LoadMetadata直接载入 proto(见上文"实现演进"),并额外做了两件事:对report.Degraded/report.Missing记录调试日志;把每个函数的原始定义(美元引用函数体,findDollarQuotedBody提取)存入funcOrigDefs,供后续 PL/pgSQL 函数体分析使用。

阶段三:核心 getQuerySpan 管线(Task 3)

这是主流程:解析 → 分类查询类型 → 分析 SELECT → 提取血缘

func (e *omniQuerySpanExtractor) getQuerySpan(ctx context.Context, stmt string) (*base.QuerySpan, error) { e.ctx = ctx // Step 1: Parse with omni. omniStmts, err := ParsePg(stmt) if err != nil { return nil, errors.Wrap(err, "failed to parse statement") } if len(omniStmts) != 1 { return nil, errors.Errorf("expected 1 statement, got %d", len(omniStmts)) } // Step 2: Extract accessed tables using omni (already migrated). accessTables, err := ExtractAccessTables(stmt) if err != nil { return nil, err } accessesMap := make(base.SourceColumnSet) for _, resource := range accessTables { accessesMap[resource] = true } // Step 3: Check for mixed system/user tables. allSystems, mixed := isMixedQuery(accessesMap) if mixed { return nil, base.MixUserSystemTablesError } // Step 4: Classify query type (already uses omni). queryType, isExplainAnalyze := classifyQueryType(omniStmts[0].AST, allSystems) if queryType != base.Select { return &base.QuerySpan{ Type: queryType, SourceColumns: base.SourceColumnSet{}, Results: []base.QuerySpanResult{}, }, nil } if isExplainAnalyze { return &base.QuerySpan{ Type: queryType, SourceColumns: accessesMap, Results: []base.QuerySpanResult{}, }, nil } // Step 5: Initialize catalog and analyze SELECT. selStmt, ok := omniStmts[0].AST.(*ast.SelectStmt) if !ok { return &base.QuerySpan{ Type: base.Select, SourceColumns: accessesMap, Results: []base.QuerySpanResult{}, }, nil } if err := e.initCatalog(); err != nil { return nil, errors.Wrap(err, "failed to init catalog") } query, err := e.cat.AnalyzeSelectStmt(selStmt) if err != nil { // Graceful degradation: return what we have. return &base.QuerySpan{ Type: base.Select, SourceColumns: accessesMap, Results: []base.QuerySpanResult{}, }, nil } // Step 6: Extract lineage from analyzed query. results := e.extractLineage(query) allSourceCols := e.extractAllSourceColumns(query) for col := range allSourceCols { accessesMap[col] = true } return &base.QuerySpan{ Type: base.Select, SourceColumns: accessesMap, Results: results, }, nil }

管线中几个关键判定:

  • 混合查询检测isMixedQuery同时命中系统表与用户表时返回base.MixUserSystemTablesError,拒绝分析。该逻辑在 access_tables.go 中实现——用户表pg_database与系统表pg_database的区别在于是否带 schema 限定(isSystemResource判断)。
  • 非 SELECT 提前返回:DML/DDL/EXPLAIN 等直接返回空结果集,只保留类型信息。
  • EXPLAIN ANALYZE:返回访问表集合但无结果列。

阶段四:列血缘 walker(Task 4,核心逻辑)

omni/pg/catalog/query_span_test.go中的概念验证 walker 移植为生产代码,输出从测试专用类型改为base.QuerySpanResult。核心映射规则:

表达式/来源解析方式
VarExpr通过RangeTable[RangeIdx]解析出 schema/table/column
RTERelation(物理表)Catalog.GetRelationByOID()Relation.Schema.Name+Relation.Name
RTESubquery递归进入Subquery.TargetList[colIdx]
RTECTE递归进入CTEList[CTEIndex].Query.TargetList[colIdx]
RTERelation+RelKind='v'(视图)递归进入Relation.AnalyzedQuery(Gap 1 修复)

与测试 walker 的关键差异:

  • 输出base.QuerySpanResult(含NameSourceColumnsIsPlainField)。
  • ColumnResource.Database使用e.defaultDatabase
  • 集合运算:合并Query.LArgQuery.RArg的血缘。

extractAllSourceColumns则行走整个QueryTargetList+JoinTree.Quals+JoinExprNode.Quals+HavingQual)收集所有被访问的列。

当前实现的增量:仓库中的extractLineage(见 query_span.go)在计划之上还做了三件事:用buildPlainFieldMask依据语法树判断IsPlainField(仅SELECT */t.*展开的列才算 plain field);isUltimatelyPlainColumn沿 CTE/子查询递归校验列是否最终落到物理表;当 catalog 把表达式折叠为常量(如json_object('id': a))丢失列引用时,回退到语法树ResTargetfigureResTargetName补名字、plpgsqlAnalyzer.extractColumnRefsFromExpr补血缘。

阶段五:集合运算(Task 5)

UNION/INTERSECT/EXCEPT会在顶层产生SetOp != SetOpNoneQuery,其TargetList中只有占位VarExpr,没有真实来源。处理方式:

  • 递归取LArg/RArg的血缘;
  • 按输出列位置合并两分支的来源列;
  • 列名取左分支(PostgreSQL 约定);
  • EXCEPT 特例:只保留左分支来源(右分支仅作过滤,不贡献输出)。

当前实现(extractSetOpLineageWithVisited)确认了这一约定:includeRight := q.SetOp != catalog.SetOpExcept && q.SetOp != catalog.SetOpExceptAll,且集合运算结果列的IsPlainField恒为false

阶段六:视图穿透血缘(Task 6)

resolveVar遇到RTERelation且底层Relation.RelKind == 'v'(视图)时,递归进入视图定义:

case catalog.RTERelation: rel := e.cat.GetRelationByOID(rte.RelOID) if rel == nil || rel.Schema == nil { return } // View through-lineage: recurse into view definition. if rel.RelKind == 'v' && rel.AnalyzedQuery != nil { if colIdx >= 0 && colIdx < len(rel.AnalyzedQuery.TargetList) { te := rel.AnalyzedQuery.TargetList[colIdx] e.walkExpr(rel.AnalyzedQuery, te.Expr, seen, result) } return } // Physical table: terminal case.

物化视图(RelKind == 'm')按同样方式处理。这样SELECT ssn FROM v的血缘能一路穿透到基表t.ssn,这正是数据脱敏(masking)所需要的"到底读了哪张表的哪一列"。测试用例TestQuerySpanResolvesAViewOverABrokenTable验证了视图上叠破损表(缺枚举类型)时血缘仍能解析到records表的列。

阶段七:系统函数作为表源(Task 7)

omni 分析器通过RTEFunction处理 FROM 子句中的系统函数。各函数的列名约定:

函数输出列
generate_series单列generate_series
generate_subscripts单列generate_subscripts
unnestN 列unnest(每个数组参数一列)
jsonb_each/json_eachkey,value
jsonb_array_elements/json_array_elementsvalue
json_to_record/jsonb_to_recordalias 子句声明的列
json_to_recordset/jsonb_to_recordsetalias 子句声明的列

计划策略:先验证 omni 是否正确设置rte.ColNames;若已正确则无需任何代码——RTEFunction的血缘 walker 天然不返回来源列(函数结果没有基表来源),列名直接取自rte.ColNames。当前实现的resolveVarRTEFunction分支进一步区分:用户自定义函数则穿透其函数体血缘,内置函数则穿透其参数的血缘(如jsonb_each(a)依赖列a)。

阶段八:用户自定义函数桥接(Task 8,复杂桥接)

BYT-9082 完成前,UDF 调用需要桥接回 ANTLR:

  1. walkExpr遇到FuncCallExpr时,通过 catalog 的UserProc注册表判断是否为用户自定义函数;
  2. SQL 语言函数:用pg.Parse()解析函数体,AnalyzeSelectStmt()分析,提取血缘;
  3. PL/pgSQL 函数:回退到既有querySpanExtractor.findFunctionDefine()逻辑。

当前实现在此基础上扩展为完整的函数体分析子系统:funcBodyCache(按函数 OID 缓存分析结果)、funcSourceColumns(函数体内发现的表级访问,合并进顶层SourceColumns)、funcPredicateColumns(函数体内 WHERE/JOIN 谓词列,独立暴露给调用方决定是否进入顶层PredicateColumns)。仓库中还有一组专门测试防止递归分析死循环(自引用函数、循环子查询、循环 CTE、循环集合运算),通过 1MB 栈的子进程验证(见 query_span_test.go 的mustNotOverflow)。

阶段九:sourceColumns 收集(Task 9)

既有 QuerySpan 的SourceColumns是一组Column字段为空的ColumnResource(表级访问),用于数据脱敏判断访问了哪些表。实现:行走Query.RangeTable,对每个RTERelation输出ColumnResource{Database, Schema, Table, Column: ""};函数体分析产生的额外来源列合并进来。

从当前源码看,表级访问已前移到管线 Step 2 的ExtractAccessTables()(非致命失败——SET等语句没有表引用),函数体列随后合并入accessesMap,因此血缘提取阶段无需重复收集。

阶段十:边界情况与错误恢复(Task 10)

  • ResourceNotFoundError:表/列不存在导致AnalyzeSelectStmt失败时,返回部分QuerySpan并设置NotFoundError
  • FunctionNotSupportedError:函数无法分析时同样返回部分结果。
  • EXPLAINEXPLAIN SELECT ...应提取内部 SELECT;EXPLAIN ANALYZE只返回访问表集合。classifyQueryType已覆盖(见 query_type.go,isExplainAnalyzeOmni通过检查 Options 中的analyzeDefElem 判断)。

当前实现的错误恢复比计划更细:解析失败时把 omni 的ParseError转成带行列位置的base.SyntaxErrorByteOffsetToRunePosition);AnalyzeSelectStmt失败时先尝试tryUserFuncTableSource(处理RETURNS TABLE函数作为表源的场景),再 fail-open 返回extractFallbackColumns的尽力而为结果,并附带UnresolvedColumnsError(通过relationHasNoSyncedColumns独立于血缘检查"元数据中该表同步了但零列"的异常)。TestGetQuerySpanNilMetadata保证从未同步过的数据库返回明确错误而非 panic。

阶段十一:清理遗留 ANTLR 代码(Task 11)

# 删除遗留提取器(3,868 行)与 ANTLR 访问表 listener(96 行) # 检查包内是否还有 ANTLR 引用 grep -r "antlr4-go/antlr" bytebase/backend/plugin/parser/pg/query_span*.go # 全量回归 go test -v -count=1 -run ^TestGetQuerySpan$ github.com/bytebase/bytebase/backend/plugin/parser/pg # Lint 与构建 golangci-lint run --allow-parallel-runners bytebase/backend/plugin/parser/pg/... go build -ldflags "-w -s" -p=16 -o ./bytebase-build/bytebase ./backend/bin/server/main.go

提交:git commit -m "refactor(pg): remove legacy ANTLR QuerySpan extractor"

注意:以上grep中使用rm删除文件的步骤属于迁移提交内容,仓库为只读,本文仅作计划还原说明,不涉及对当前仓库的修改。


四、任务汇总与风险评估

Task描述预估行数变化风险
1提取器脚手架 + 入口接线+80
2catalog 元数据加载+50
3核心 getQuerySpan 管线+80
4列血缘 walker+200高(核心逻辑)
5集合运算+40
6视图穿透血缘+20
7系统函数+20(验证为主)
8UDF 桥接回 ANTLR+100高(复杂桥接)
9sourceColumns 收集+30
10边界情况 + 错误恢复+50
11删除遗留代码-3,964低(纯删除)

净结果:约 +670 行、-3,964 行,净移除约3,300 行


五、测试策略:77 个 YAML 用例作为验收标准

所有既有 YAML 用例就是本次迁移的验收标准,测试运行器 query_span_test.go 的TestGetQuerySpan本身不改动——它调用入口GetQuerySpan(),遍历 test-data/query_span.yaml(计划口径 46 个用例)与 test-data/query_type.yaml(31 个用例),将结果序列化为 YAML 与 golden 数据比对(result.ToYaml()tc.QuerySpan逐字段相等),共 77 个用例。

每个用例的结构:

- description: 用例描述 statement: SELECT ... # 被测 SQL defaultDatabase: db # 默认数据库 metadata: '...' # protojson 编码的 DatabaseSchemaMetadata querySpan: # 期望的 golden 结果 type: SELECT sourceColumns: [...] results: [...]

测试运行命令(每个任务完成后执行,不允许累积失败):

go test -v -count=1 -run ^TestGetQuerySpan$ github.com/bytebase/bytebase/backend/plugin/parser/pg

此外,测试套件还包含一批计划之外的健壮性用例,集中体现了"降级但不崩溃"的工程取向:

  • 坏引号标识符('weird'table)不阻塞同库其他表查询;
  • 引用未声明枚举类型的表以 stand-in 安装,列名仍可解析;
  • 同一 schema 中坏表不影响健康表的血缘(爆炸半径控制);
  • 分区表血缘解析到分区本身(orders_2024),供脱敏将分区解析回父表;
  • WITH ORDINALITY的 LATERAL 函数、jsonb_path_query_array等复杂 JSONB 表达式在 CTE 链路中血缘保持完整。

六、总结

PostgreSQL QuerySpan 迁移到 omni 是一次教科书式的"以语义分析替代语法遍历"重构:让解析器替我们完成名称解析、作用域与星号展开,血缘提取器只负责行走已经"懂语义"的树。这不仅把 3,868 行的手工 ANTLR 遍历压缩到数百行,更重要的是把正确性责任移交给了可复用的语义分析基础设施——后续 MySQL、MSSQL 的同类迁移(仓库中已有 2026-04-23-mssql-query-span-omni-migration.md、2026-04-27-mysql-query-span-omni-migration.md 计划)可以复用同一套模式。

对希望复刻本次经验的团队,核心可借鉴点有三:

  1. 先建"概念验证测试"再写生产代码:把血缘 walker 先放在 omni 的测试里验证正确性,再移植到生产侧适配base.QuerySpanResult
  2. 以 golden YAML 全量回归作为安全网:77 个既有用例不动、测试运行器不动,任何一步的语义漂移都会立刻暴露;
  3. fail-open 优于 fail-hard:分析失败时返回部分结果(访问表 + 尽力而为的列),并设置明确的错误标志,让上层(脱敏、审批、SQL 审核)自行决策,而不是让一次无法分析打断整个功能链路。

【免费下载链接】bytebaseDatabase governance built for humans and agents — controlling changes and access across every major database.项目地址: https://gitcode.com/GitHub_Trending/by/bytebase

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

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

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

立即咨询