构建以数据为中心的软件工厂
十年前,管道是一个GUI配置的ETL作业、一个有人拥有root权限的机器上的定时脚本,以及一个描述它大致做什么的wiki页面。没有什么东西是智能体可以抓住的——没有仓库、没有测试套件、没有编译器、没有CI。
管道不再是那样。转换是dbt模型。编排是Python DAG定义。摄取是声明式配置。提取是Spark代码。所有这些都存在于版本控制中,作为差异到达,通过CI运行,并像任何其他软件更改一样进行评审——附加相同的规范:代码检查、单元测试、代码评审、分阶段环境、语义版本控制。
这种转变是本文所有内容的前提。因为工件是存放在具有代码检查器、编译器和测试框架的仓库中的代码,我们可以做一件使自主性变得可容忍的事情:在智能体工作落地之前机械地检查它。 您无法在pytest断言后面设置拖放GUI映射的关卡。您绝对可以设置dbt模型的关卡。
因此,软件工厂的想法——将智能体输出视为具有质量关卡的工业过程,而不是您目视的聊天记录——转移到数据工程,正是因为数据工程首先采用了软件工程的实践。
这是我们开始测试的赌注。
因此,我们构建了一个概念验证,其中一个AI智能体自主检测上游源中的模式漂移,跨现代数据工程堆栈(dbt、Airflow、Spark、Snowflake)生成代码更改,使用充当"基本事实预言机"的确定性测试工具验证每个工件,并监控生成的管道直到它成功或失败——智能地重试并在无法恢复时升级。核心赌注:在确定性后面设置AI关卡。 LLM会产生幻觉;测试工具不会。LLM生成的每个字节都必须通过机器可检查的关卡,才能被允许接触真实基础设施。
注意:本文中的实现细节是概括性的。命名、模式字段和层约定是说明性的,而不是从任何特定环境中提取的。
1、问题:模式漂移是千刀万剐
典型的摄取层从异构源——文档存储、关系数据库和平面文件——拉取到勋章架构(登陆→原始→精选),该架构在Airflow、dbt和Spark上针对云仓库运行。
源模式从不稳定。列被添加。数据类型更改。列被弃用。单个上游模式更改可以级联通过摄取SQL模板、摄取配置、转换配置、dbt模型,有时还有多个编排定义和提取作业。
手动处理每次漂移意味着:读取元数据差异,计算影响半径,跨多个仓库编辑文件,推送,触发管道,然后监视日志直到它成功或出现细微故障——然后调试、重新生成和重试。
我们问了一个难题:我们可以端到端地自动化这个吗?
2、设计
我们构建了一个三层智能体系统,如下所示
1. 工程智能体 — "构建者。"它从目录API(如Openmetadata)读取模式元数据,分类发生了什么类型的更改,并生成相应的工件。分类很重要,因为影响半径因更改类别而有很大差异:
两件事使这种分类法有用而不是学术性的。
- 首先,风险不是均匀的 — 添加性更改是低风险,类型更改是危险的,删除需要故意弃用而不是删除。
- 其次,触及的仓库数量随类别增长,因此提前知道类别告诉智能体它实际上需要加载多少上下文。只有在引入全新类型源系统的情况下才需要新的提取作业,而不仅仅是现有系统中的新命名空间。
2. 质量智能体 — "守护者。"在新表落地后,它使数据质量声明与模式现实保持同步。关键的是,它写入现有DQ框架(验证配置加上dbt测试定义)而不是构建一个并行框架——因此生成的规则自动参与现有的隔离表、严重性分类和警报,零新基础设施。
3. 编排器 — "管理者。"它排序所有内容,拥有重试循环(上限两次重试),决定何时运行质量检查,启动监控监听器监视管道链,并在无法恢复时升级给人工。
4. 持久内存和来源 一个不能记住上次做了什么的智能体无法检测自上次以来发生了什么变化。持久性使"检查此表"成为有意义的指令,而不是完全重新推导的请求。每个智能体保持自己的状态文件,每个表一个键。在该键内有三个独立的关注点:
{
"relational/orders": {
"last_applied_class": "additive",
"last_applied_version": 1.1,
"last_applied_at": "2026–08–04T00:00:00Z",
"downstream_impact": {
"raw_table": "<db>.<raw_schema>.<table>",
"curated_table": "<db>.<curated_schema>.<table>",
"model": "models/curated/<…>.sql",
"orchestration": "<pipeline_id>",
"affected_columns_by_class": {
"additive": [
"XX",
"YY"
]
}
},
"commit_history": [
{
"class": "additive",
"version": 1.1,
"repos": {
"transform": {
"hash": "4810daf",
"message": "feat: add columns"
}
},
"risk": "LOW",
"columns_changed": [
"XX",
"YY"
]
}
]
}
}
游标是运行时唯一需要的部分 last_applied_* 是游标。它是最小的契约——自主性工作必须持久保存的一件事,因为更改类别形成顺序链。
有了游标,一个裸检查
解析到要差异的特定模式快照对。没有它,"自上次以来发生了什么变化"没有参照物。剥离每个其他字段,自动检测仍然有效。
新实体和新命名空间类位于链之外,独立于游标。
Commit_history 是日志——仅追加,last_applied_version 实际上是指向它的指针。每次运行推进游标并附加一个条目。
我们利用了Snowflake的Cortex Code来执行Snowflake特定任务以及其余数据工程管道。
3、预言机:使这安全的部分
这是我们构建的模式。我们不信任LLM了解我们自己的约定。 这些作为确定性Python测试工具存在——一个编码以下内容的"契约"预言机:
- 允许的测试类型和验证规则类型的精确集合
- 仓库目录约定和文件命名规则
- 每个转换模型必须遵循的结构化CTE模式
- DDL操作和元数据列要求
- 针对实际仓库的命名正确性
在任何生成的工件被推送之前,智能体运行一个验证子技能——对智能体刚刚编辑的实际工作树执行这些pytest工具,而不是对模拟。这是"智能体说它看起来正确"和"机器验证它是正确的"之间的区别。
工具优雅地降级——如果仓库克隆或实时基础设施不存在,它们跳过并显示消息而不是失败。这保持了检查的诚实性:绿色运行意味着基础设施在那里并且工件通过了。
4、预言机实践:工具脚本实际做什么
"确定性工具"听起来很抽象。这是它具体的样子。我们编写了一组小型Python模块,编码了LLM绝不允许猜测的事实。下面的代码片段是我们构建内容的说明。
4.1 漂移分类 — 纯Python,零LLM
给定来自目录的先前和当前模式快照,这返回更改类别、风险评级和确切的列更改:
# classify.py - deterministic, no LLM in the hot path
@dataclass
class DriftResult:
change_class: str
risk: str
added: list[str] = []
removed: list[str] = []
type_changes: list[tuple[str, str, str]] = [] # (name, old_type, new_type)
def classify_drift(prev_schema: dict, curr_schema: dict) -> DriftResult:
prev_cols = {c["name"]: c for c in prev_schema["columns"]}
curr_cols = {c["name"]: c for c in curr_schema["columns"]}
added = sorted(set(curr_cols) - set(prev_cols))
removed = sorted(set(prev_cols) - set(curr_cols))
type_changes = sorted(
(n, prev_cols[n]["dataType"], curr_cols[n]["dataType"])
for n in set(prev_cols) & set(curr_cols)
if prev_cols[n]["dataType"] != curr_cols[n]["dataType"]
)
# Priority matters: a type change outranks an addition, because it is
# the change most likely to corrupt downstream data silently.
if type_changes: return DriftResult("type_change", "MEDIUM-HIGH", added, removed, type_changes)
if added and not removed: return DriftResult("additive", "LOW", added, removed, type_changes)
if removed: return DriftResult("removal", "HIGH", added, removed, type_changes)
return DriftResult("no_change", "NONE", added, removed, type_changes)
智能体委托给这个而不是通过散文推理差异——更快,而且它不能幻觉更改类别。注意显式的优先级排序:当表既获得列又更改类型时,类型更改获胜,因为这是悄悄产生错误数字的那个。
4.2 契约预言机 — 唯一的事实来源
此模块定义每个工件必须满足的不变量。最重要的:一组管道管理的元数据列,任何更改类别都不允许删除、重命名或重新键入。这些是平台本身依赖的列,用于沿袭、去重和增量逻辑——如果智能体触及它们,管道会以难以跟踪的方式中断。
# contracts.py - the "sacred" invariants (illustrative names)
# Pipeline-managed metadata columns: NEVER dropped, renamed, or retyped.
METADATA_COLUMNS = frozenset({
"RECORD_KEY", # surrogate key
"RAW_PAYLOAD", # original source document
"SOURCE_PRIMARY_KEY",
"SOURCE_MODIFIED_AT",
"SOURCE_ROW_HASH", # change detection
"EXTRACTED_AT",
"SOURCE_SYSTEM",
"SOURCE_OBJECT",
"SOURCE_FILE_PATH",
"LOADED_AT",
"SNAPSHOT_DATE",
})
4.3 模式检查 — 断言模型匹配内部风格,而不是通用dbt
这些不是读取代码,而是对生成的SQL进行断言。每个转换模型必须携带所需的config()块、正确的unique_key、预期的CTE结构,以及从去重阶段的最终选择:
# model_checks.py - structural pattern enforcement
REQUIRED_CONFIG = {"materialized": "incremental", "incremental_strategy": "merge"}
DEDUP_CTE = "deduplicated_batch" # house-style CTE name
def check_config_block(model_sql: str) -> list[str]:
m = re.search(r"\{\{\s*config\((.*?)\)\s*\}\}", model_sql, re.IGNORECASE | re.DOTALL)
if not m:
return ["missing required config() block"]
block = m.group(1)
violations = []
for key, expected in REQUIRED_CONFIG.items():
if not re.search(rf"{key}\s*=\s*['\"]{re.escape(expected)}['\"]", block, re.IGNORECASE):
violations.append(f"config() missing/incorrect {key}='{expected}'")
return violations
def check_dedup_stage(model_sql: str) -> list[str]:
if not re.search(rf"\b{DEDUP_CTE}\b", model_sql, re.IGNORECASE):
return [f"model must build a '{DEDUP_CTE}' CTE"]
return []
一个微妙但重要的细节:预言机是从仓库实际执行的内容转录的,而不是从文档中的理想化模板。这两个随时间漂移,当它们漂移时,仓库是事实。工具还携带预先存在违规的显式豁免列表,因此新运行不会因未引入的漂移而受到指责。
4.4 基础 — 智能体不猜测东西在哪里
仓库位置、目标分支和目录端点由代码解决,具有显式回退——大声跳过,从不静默通过:
# repos.py - resolve real repo clones, or skip loudly
REPO_DIRS = {"transform": "analytics-transform", "orchestrate": "analytics-orchestration"}
REPOS_ROOT = Path(os.environ.get("REPOS_ROOT", default_root))
BRANCH = os.environ.get("REPO_BRANCH", "development")
def repo_available(kind: str) -> tuple[bool, str]:
path = REPOS_ROOT / REPO_DIRS[kind]
if not path.is_dir() or not (path / ".git").exists():
return False, f"{kind} repo not found at {path} - clone it or set REPOS_ROOT"
return True, "ok"
5. 将其连接到智能体生命周期
验证是生成和推送之间的一等步骤。验证技能将更改类别映射到精确的pytest选择——它从不"目视"代码:
# validation skill: change class -> pytest selection
selections = {
"additive | type_change | removal": 'pytest tests/test_column_drift.py -k "<class> and <source>"',
"new_entity": 'pytest tests/test_new_entity.py -k "<source>"',
"new_namespace": 'pytest tests/test_new_namespace.py -k "<source>"',
}
循环现在是:在Python中分类 → 在LLM中生成 → 使用pytest针对真实工作树验证 → 推送。
5.1 闭合循环:监视管道,而不仅仅是推送
生成看起来正确的代码与交付工作管道不同。因此编排器启动一个后台监控监听器,监视数据集触发的管道链并闭合反馈循环。
监听器不仅报告通过/失败——它分类失败,因为正确的响应完全不同:
- 瞬态基础设施 — 完全重新生成并让新运行触发
- 模式或代码错误 — 将错误反馈给工程智能体作为上下文,以便它在考虑失败的情况下重新生成受影响的工件
- 配置错误 — 不要重试;立即升级,因为重试配置错误的监视器只会消耗另一个长时间超时
第三种情况是团队通常搞错的。智能体无法通过重新生成修复的失败必须立即退出循环,而不是消耗重试预算。
5.2 构建MR模板
监控闭合了智能体和管道之间的循环。但还有第二个循环,这是我们低估的:智能体和必须处理其输出的人类之间的循环。
推送代码的智能体必须回答一个人在做时永远不会出现的问题:评审者如何验证他们没有编写的工作?
我们的答案是一个生成的评审工件——每个仓库一个合并请求描述。四个部分承担了重量:
1)上游更改,以差异表示法。 一个列标记为+添加、~修改、-删除的表,带有类型和可空性。评审者看到原因,而不仅仅是结果。
2)机器自己的证据,作为复选框。
这是工具输出显示到评审中。评审者不被要求信任智能体——他们被显示它清除了哪些确定性关卡。关键的是这些是报告的,从不声称的:智能体不能勾选工具未返回的框。
3)一个显式的影响部分,回答评审者否则必须自己推导的问题:
4)交叉链接到同级请求。 一个上游模式更改可能触及三个仓库,提交不能跨它们原子化。每个描述链接其他描述。没有这个,评审者看到一个更改的三分之一,无法知道其余部分存在——这是评审多仓库自动化输出时最令人困惑的事情。
5.3 我们学到了什么
- LLM在一般性上训练;您的仓库是一个特定的事实集。 将每个约定、命名规则和结构要求移动到确定性代码中。将工具视为契约,将智能体视为必须满足它的解释器。
- 最便宜的token是从未花费在LLM上的那个。 漂移分类、配置解析和验证都是确定性的——将它们委托给Python而不是推理它们。更快,而且它从不产生幻觉。
- 排序是架构,而不是约定。 质量检查运行多少次、管道何时运行、以及哪个智能体在失败时重试是具有真实影响半径后果的决策。让编排器明确拥有它们:质量失败绝不应错误地归因于工程智能体,配置失败绝不应触发无意义的重新生成和重试循环。
- "代码正确"和"管道工作"是两个不同的声明。 验证涵盖第一个;带有失败分类的监控循环涵盖第二个。在您让智能体推送到共享分支之前,两者都是必需的。
- 大声跳过,从不静默。 每个无法运行的工具检查都报告原因。当基础设施不存在时静默通过的测试套件比没有测试套件更糟,因为它制造了虚假的信心。
6、结束语
这不是一个完成的剧本。我们仍在构建,接下来的很多内容是我们先做错后学到的。
这个概念验证证明了循环有效。但"它有效"与"它以正确的成本运行"不同。今年的行业对话已从 "智能体能做到吗?"转变为 "我们能在不增加账单的情况下扩展智能体吗?"
第2部分将该镜头应用于此系统。我们借鉴了Uber关于高效运行软件工厂的公共工程手册——他们的前提是智能体成本分解为一个方程,您可以独立攻击其项:
> _cost ≈ users × sessions × turns × requests × tokens × price
您不只是协商更低的token价格。您消除零值token——那些花费在LLM根本不应该做的工作上的token。对于这个智能体工厂,这映射到八个具体的改进:
- 代码模式分类和路由 — 使用Python CLI分类漂移并返回路由卡,而不是让LLM用散文重新推导差异。分类从不触及LLM token。
- 静态上下文索引 — 从每个表到其沿袭和预期工件的预构建映射,用单个查找替换仓库grep。大型商店维护的大型上下文图的轻量级替代品。
- 精简编排提示 — 今天编排器在每次运行时加载完整的智能体技能文件(超过一千行),即使单列更改。用只加载更改类别需要的内容的路由卡替换它。
- 更改类别门控参考加载 — 只读取更改类别实际接触的模式文件,而不是每次读取每个模式文件。
- 结果指标,仓库原生 — 持续时间、失败率和每个更改类别的重试作为仓库视图,因此"托管智能体单位经济"不再是氛围而成为表。
- 小伤反馈循环 — 持久化分类的管道失败,以便下一个技能迭代由真实失败而不是记忆驱动。
- 分层成本控制和操作员默认值 — 明确哪个表面是托管的与交互的,并记录减少每请求token的上下文压缩和推理努力默认值。
- 评估矩阵 — 重用现有pytest工具作为基准,因此任何模型或技能更改在发布前都根据相同的栏进行评判。
贯穿始终的主题: 最便宜的token是从未发送到LLM的那个。敬请关注这段旅程。
原文链接:Building a Data-Focused Software Factory
汇智网翻译整理,转载请标明出处