ADR-006-data-rule-runtime.md 42 KB

ADR-006:AI 驱动的数据规则生成、编译与运行时执行

  • 状态:Proposed
  • 日期:2026-07-21
  • 代码基线:234df2emaster
  • 范围:数据标准与数据流程中的规则定义、数据生产线装配、数据工厂投产、执行、审计和迁移

1. 结论摘要

DataOps 应把自然语言作为数据规则管理的一级入口。用户用中文描述约束、清洗和 加工要求后,平台内置的 AI Rule Authoring Agent 自动补充字段 Schema、数据画像、 已有规则和目标数据源方言等上下文,把意图转换为版本化 RuleSpec;当声明式算子 不足以表达规则时,生成受治理的 GeneratedCodeSpec + 代码 + 测试。平台随后自动 完成确定性校验、编译、样本试运行、风险判定、发布和 DataFlow 绑定,避免用户手写 规则表达式、SQL、Python 或流程胶水代码。

自然语言是不可变的 authoring source(意图源),经过验证的 RuleSpec、表达式或 签名代码产物是 runtime source(执行源);两者通过版本和 Hash 关联。大模型参与 规则理解、生成、解释、修复和影响分析,但不进入逐行数据处理热路径。低风险规则可 自动发布和绑定;只有存在语义歧义、低置信度或高风险写入时才要求人工介入。

这套能力必须同时覆盖三个产品域,并共享同一规则内核,不能只实现在“数据流程”:

  1. 数据标准定义字段、数据域、模型或数据产品必须满足的业务含义、格式、质量和 合规约束;AI 把标准条款生成一个或多个可执行 RuleVersion,并与 StandardVersion 固定绑定。
  2. 数据流程(数据生产线)按顺序组装输入输出 BusinessDomain、已发布的数据标准 和通用/流程专用数据规则,形成不可变 DataFlowVersion。DataFlow 是生产线的治理 定义,不等同于某个调度引擎的工作流 JSON。
  3. 数据工厂把已发布的数据生产线绑定到具体环境、数据源、算力、调度和运行策略, 编译为 WorkflowSpec/Kestra Flow,完成 canary、投产、监控和回滚;产出是可交付的 数据产品。

规则执行采用以下组合架构:

  1. CEL 表达简单条件、过滤和派生表达式;它是非图灵完备、无副作用、适合嵌入 应用的开源表达式语言。
  2. SQLGlot 负责 SQL AST、方言校验和安全下推;同库、可下推的规则优先在源库 执行,避免搬运数据。
  3. Polars Lazy/Streaming 负责跨库、文件和不适合 SQL 下推的列式批处理。
  4. Great Expectations(可选) 作为复杂数据质量规则适配器,不作为转换主引擎。
  5. Kestra 继续负责调度、重试和运行状态机;DataOps Runner 新增 rule.apply / quality.check 执行器,继续掌握凭据、连接池、幂等和审计。
  6. AI Rule Authoring Agent 复用现有 AI 调度规划的 Schema 约束、确定性校验、 canary 和审计模式,形成“自然语言 -> 可执行规则版本”的自动闭环。

首期不引入 Drools、dbt Core 或 Apache Beam。它们分别更适合复杂业务决策、 SQL 工程项目和大规模分布式批流处理,作为当前规则运行内核都会增加不必要的 运行时和治理边界。

2. 当前实现的事实基线

2.1 当前链路

flowchart LR
    STANDARD["数据标准页面\n描述 + input/output + 操作代码"] --> SCODE["LLM 直接生成 Python 文本"]
    SCODE --> SNODE["Neo4j data_standard\n保存 code 属性"]

    UI["数据流程页面\n自由文本 rule"] --> DF["Neo4j DataFlow\nscript_requirement JSON 字符串"]
    DF --> TASK["PostgreSQL task_list\nMarkdown 任务说明"]
    TASK --> CODE["外部流程生成 Python 脚本"]
    DF --> N8N["本地生成 n8n Workflow JSON"]
    N8N --> FACTORY["数据工厂\n激活、触发、执行记录"]
    FACTORY --> SSH["SSH 执行 Python 脚本"]

    SPEC["WorkflowSpec"] --> KESTRA["Kestra DAG"]
    KESTRA --> RUNNER["DataOps Runner"]
    RUNNER --> SQL["受控 SQL / 受限 Python / HTTP"]

代码证据:

  • frontend/src/views/dataGovernance/dataStandard/components/edit.vue 已提供自然语言 describe、输入/输出参数和“代码生成”,并把模型返回文本直接放入必填 codeapp/api/data_interface/routes.py 将其作为 Neo4j data_standard 节点属性保存。
  • app/core/llm/code_generation.py 直接要求模型生成 Python 函数,没有结构化输出、 版本、测试、制品签名和执行绑定;当前接口的参数键名也不一致,说明它还只是原型。
  • frontend/src/views/dataGovernance/dataProcess/components/edit.vue 使用一个 v-textarea 接收 changeObj.rule,提交时把自由文本与源/目标 BusinessDomain 一起包装成 script_requirement
  • app/core/data_flow/dataflows.py:238-288script_requirement 序列化到 Neo4j;:383-610 将自由文本拼成 Markdown 任务说明,写入 task_list,再生成 n8n Workflow JSON。
  • app/core/data_flow/dataflows.py:660-708 生成 SSH 命令来运行约定路径下的 Python 脚本。规则本身没有经过结构化校验、编译或版本绑定。
  • app/api/data_factory/routes.py 已承担工作流查询、激活/停用、触发和执行记录等投产 操作,但接收的是引擎工作流,尚未把“已发布数据生产线”作为投产对象。
  • app/core/data_flow/dataflows.py:1132-1244 编辑 DataFlow 时只更新 Neo4j 属性; 不会生成新的不可变规则版本,也不会同步重建 task_list、脚本或执行定义。
  • app/core/data_processing/data_validator.pydata_cleaner.py 已有 Pandas 清洗/校验工具,但仓库其他运行代码没有引用它们,尚未进入生产执行链路。

2.2 已有可复用基础

  • app/core/orchestration/spec.py 已有封闭字段、DAG、秘密字段、幂等、重试和回填 校验,并对规范内容生成稳定哈希。
  • app/core/orchestration/compilers/kestra.py 已能把 WorkflowSpec 确定性编译成 禁用状态的 Kestra Flow,并把每个节点包装为带短期任务令牌的 Runner 请求。
  • app/runner/nodes.py 已实现只读 SQL、单条参数化 DML、写入幂等约束、HTTP allowlist 和资源受限的预注册 Python Handler。
  • app/runner/auth.pyrunner_task_executions 已实现节点摘要绑定、短期令牌、 单次消费和提交结果记录。
  • 数据源凭据、连接池、熔断器和事务边界已经集中在 Runner 可复用的 DataSourceConnectionManager 中。
  • MinIO 已是平台对象存储,可承载跨节点的 Parquet 中间产物和隔离样本。
  • app/core/data_service/data_product_service.py 已能调用大模型把自然语言需求提取为 output_domain/key_fields/processing_logicapp/core/llm/code_generation.py 也已有 基础 Python 代码生成入口;但当前输出是宽松 JSON/文本,没有规则 Schema、测试、 沙箱、版本和执行绑定。
  • app/core/orchestration/agent 已实现 AI 调度规划的封闭 JSON Schema、授权资源校验、 disabled 部署、canary、内容 Hash 和模型/Prompt/上下文审计。这套控制器模式可直接 复用于规则生成,而不是另起一套不可审计的 Agent 链路。

2.3 必须补齐的缺口

缺口 当前影响 目标状态
自然语言未进入受控生成闭环 只能成为任务说明,仍需人或外部流程写脚本 AI 自动生成、验证、测试和发布 RuleVersion
创建和编辑链路不一致 编辑规则不会产生新执行版本 所有变更生成草稿版本,发布后才可执行
无规则到 WorkflowSpec 的编译层 Kestra/Runner 不知道规则语义 RuleSpec -> ExecutionPlan -> rule.apply
WorkflowSpec 声明与 Runner 能力不一致 声明了 8 类节点,Runner 实际注册 4 类;Python Handler 为空 发布前做“已注册且可执行”能力校验
SQL 查询先 .all() 再截断 1000 行 不能处理真实批量数据,且可能占用大量内存 服务端游标/分页、下推 SQL或 Parquet 批次
节点间不传递上游结果 当前 Kestra DAG 只传运行参数和令牌 原子 rule.apply 或 MinIO artifact_ref
无拒绝/隔离输出 脏数据只能使整个任务失败或被忽略 pass/reject/quarantine 多结果通道
CRUD 缺少统一规则级权限 规则可能绕过发布治理 rules:read/edit/publish/execute 后端强制权限
任意脚本历史路径 可重复性、安全性和供应链风险高 只运行结构化规则或已审核、签名的插件
流程选择规则后仍需手工接线 规则与流程执行定义可能漂移 自动解析版本、编译产物并注入 rule.apply
数据标准和数据流程各自生成代码 同一约束可能出现两份实现,无法统一修复和统计 共用 RuleVersion、编译器、制品和运行指标
DataFlow 与引擎工作流边界模糊 治理定义、生产线版本和投产状态相互污染 DataFlowVersion 装配,Data Factory 部署
生产线只绑定规则、不绑定标准 标准条款可能在实际生产中未执行 standard.enforce 固定 StandardVersion 并展开规则

3. 开源技术调研

3.1 评价维度

候选方案按以下维度评估:真实数据变换能力、规则可版本化/可审计、Python/Flask 适配、批量性能、多数据库方言、安全隔离、运维复杂度和许可证。

3.2 候选方案对比

技术路线 适合做什么 优点 主要限制 许可证 本项目定位
CEL 条件、过滤、简单派生表达式 非图灵完备、无副作用、可预编译、可嵌入;官方说明适合谓词和简单变换 不是批处理引擎;Python 实现需做兼容性 PoC Apache-2.0 采用:表达式层
SQLGlot SQL AST、方言转换、静态分析 Python 原生,支持 PostgreSQL/MySQL 等方言,可拒绝不支持语义 不是高性能执行引擎;跨方言转换必须严格失败而非 best-effort MIT 采用:SQL 编译层
Polars 跨库/文件的列式批量转换 Lazy 优化、谓词/投影下推、流式执行、Python 接入轻 单机为主;不是集群批流平台 MIT 采用:批处理层
Great Expectations Core 数据质量断言、Suite、验证结果 声明式质量规则、批次与验证结果模型完整 主要做验证,不负责数据转换;引入完整 DataContext 有额外成本 Apache-2.0 可选:质量适配器
Drools / Apache KIE 事实匹配、决策表、DMN、CEP 复杂规则冲突和推理能力成熟 JVM 服务和 KIE 资产体系较重;不擅长列式 ETL Apache-2.0(Incubating) 暂不采用;复杂决策场景再评估
dbt Core 仓库内 SQL 模型、依赖、测试 SQL 工程化和模型 DAG 成熟 偏项目/文件/CLI 工作流;难直接承载 UI 中的细粒度动态规则;当前主分支 v2 仍为 alpha Apache-2.0 不作为规则内核;可作为外部适配器
Apache Beam 大规模批流统一、分布式转换 批流统一,Runner 可替换,PTransform 模型完整 引入 Runner/集群/序列化/IO 体系,当前规模下运维成本高 Apache-2.0 后续规模化执行适配器

补充说明:

  • CEL 官方规范强调线性时间、无修改、非图灵完备;Python 侧可对 cloud-custodian/cel-python 做兼容性 PoC。该实现是 Apache-2.0,但其类型预检查能力与 Go/C++ 实现存在 差异,因此首期必须用 DataOps 自己的输入/输出 Schema 做二次校验。
  • SQLGlot 官方支持把不兼容转换设置为 RAISE。DataOps 必须使用严格模式并在 发布阶段绑定源/目标方言,不能使用默认 best-effort 输出。
  • Polars 官方 Lazy API 支持查询优化;Streaming 可以分批处理超过内存的数据, 但部分算子会回退到内存模式,因此编译计划必须记录算子是否可流式执行。
  • Great Expectations 的 Expectation/Suite/Validation Result 很适合质量闸门, 但“清洗/加工”和“质量校验”应保持为两个明确的执行语义。

3.3 大模型与结构化生成方案

保留现有 OpenAI-compatible 模型接口,不让规则控制面绑定某一家模型。生产建议优先 评估 vLLM + Qwen3 系列的私有化组合:vLLM 是 Apache-2.0,提供 OpenAI-compatible 服务和 JSON Schema/grammar 结构化输出;Qwen3 官方仓库声明开放权重采用 Apache-2.0,并提供中文、工具调用和代码相关模型。这样可复用当前 DataOps 的模型 客户端,同时把元数据和规则意图留在内网。具体模型与量化版本必须用企业规则样本做 准确率、歧义识别率、生成代码测试通过率、延迟和显存基准后再固定,不能只按通用榜单 选型;每个模型制品仍需单独核验权重许可证。

云端大模型可以作为可配置的高质量或容灾后端,但必须经过脱敏与出域策略。无论采用 本地还是云端模型,控制器都使用同一 RuleCandidate JSON Schema、Pydantic 二次 校验、温度 0、固定 Prompt 版本和回归集。vLLM 的 constrained decoding 只能保证 输出“形状正确”,不能保证字段含义和数据处理语义正确,因此不可替代编译、测试和 风险策略。

4. 推荐目标架构

4.1 数据标准、数据生产线与数据工厂的领域边界

flowchart LR
    subgraph GOVERN["数据治理与设计态"]
        DS["数据标准 StandardVersion\n业务口径、质量与合规条款"]
        DR["数据规则 RuleVersion\n清洗、转换、校验实现"]
        SRB["StandardRuleBinding\n标准条款固定到规则版本"]
        DS --> SRB --> DR
        DF["DataFlowVersion / 数据生产线\n阶段、依赖、输入输出"]
        DS -->|"standard.enforce"| DF
        DR -->|"rule.apply / quality.check"| DF
    end

    subgraph FACTORY["数据工厂与投产态"]
        DEPLOY["DataFlowDeployment\n环境、数据源、调度、资源、策略"]
        WF["WorkflowSpec + ExecutionPlans\n不可变生产线部署包"]
        K["Kestra Flow\n调度与状态机"]
        R["DataOps Runner\n确定性执行"]
        DEPLOY --> WF --> K --> R
    end

    DF -->|"发布生产线版本"| DEPLOY
    R --> PRODUCT["数据产品\n数据 + 质量证据 + 血缘"]

可以用实物工厂类比,但系统对象必须保持清晰:

实物生产概念 DataOps 对象 主要职责
产品/质量标准 DataStandardVersion 定义数据应当满足什么业务、格式、质量和合规要求
工艺规则/工位作业 RuleVersion 定义每一步如何校验、清洗、转换或路由异常数据
生产线/工艺路线 DataFlowVersion 按阶段装配标准和规则,声明输入、输出、依赖和质量闸门
工厂、设备和班次 DataFlowDeployment + Data Factory 绑定环境、数据源、算力、调度、容量和运维策略并投产
在制品/成品 artifact / DataProduct 中间数据制品以及带质量和血缘证明的数据产品

关键约束:

  • DataStandardVersionRuleVersion 是独立可复用资产。一个标准可拆为多个规则; 同一规则也可被多个标准或流程复用,二者不能把描述和代码互相复制。
  • DataFlow 在产品语义上就是数据生产线;建议保留现有 DataFlow 稳定 UID,新增 不可变 DataFlowVersion 和组件绑定,不再创建另一个重复的“ProductionLine”实体。
  • DataFlowVersion 的组件可以是 standard.enforcerule.applyquality.check 和 输入/输出节点。标准组件固定 StandardVersion,发布时展开并固定其 RuleVersion; 因此标准后续升级不会悄悄改变已投产生产线。
  • 数据工厂只投产已发布的 DataFlowVersion,不重新理解自然语言或生成业务规则。 它负责部署参数、编译、canary、激活、停用、回滚和运行监控,保持设计态与投产态分离。
  • 数据产品不是只有目标表;还应关联生产它的 DataFlowVersion/Deployment、规则版本、 输入版本、质量结果、血缘和运行批次,形成可验证的“产品合格证”。

4.2 AI 原生规则生命周期与责任边界

flowchart LR
    STDNL["数据标准自然语言\n标准条款与合规要求"] --> CONTEXT["上下文解析\nSchema、画像、标准/规则目录、方言"]
    FLOWNL["数据流程自然语言\n加工步骤与质量闸门"] --> CONTEXT
    CONTEXT --> AGENT["AI Rule Authoring Agent\n约束抽取、RuleSpec/代码生成"]
    AGENT --> VALIDATE["确定性控制器\nSchema、类型、安全和权限校验"]
    VALIDATE --> COMPILE["Rule Compiler + Artifact Builder"]
    COMPILE --> SQLPLAN["SQLGlot AST\n数据库下推"]
    COMPILE --> POLARSPLAN["Polars Lazy Plan\n跨库/文件批处理"]
    COMPILE --> QUALITY["内建质量检查 / GX Adapter"]
    COMPILE --> CODE["签名 Python Artifact\n受限依赖与沙箱"]
    COMPILE --> PLAN["不可变 ExecutionPlan + Hash"]
    PLAN --> TEST["自动测试与样本试运行\n失败时有界 AI 修复"]
    TEST --> POLICY{"风险与置信度策略"}
    POLICY -->|"低/中风险且通过"| PUBLISH["自动发布不可变 RuleVersion"]
    POLICY -->|"歧义、低置信度或高风险"| REVIEW["请求澄清或人工审批"]
    REVIEW --> PUBLISH

    PUBLISH --> RV["已发布 RuleVersion"]
    RV --> STANDARD["StandardVersion\n固定标准规则绑定"]
    RV --> FLOW["DataFlowVersion\n装配规则"]
    STANDARD --> FLOW
    FLOW --> RESOLVE["自动展开标准并固定兼容版本"]
    PLAN --> RESOLVE
    RESOLVE --> DEPLOY["Data Factory Deployment\n环境、资源、调度"]
    DEPLOY --> WFSPEC["自动生成 WorkflowSpec\nstandard.enforce / rule.apply"]
    WFSPEC --> KESTRA["Kestra\n调度、重试、状态机"]
    KESTRA --> RUNNER["DataOps Runner"]
    RUNNER --> POOL["现有数据源连接池"]
    RUNNER --> MINIO["MinIO Parquet artifact_ref"]
    RUNNER --> RESULT["规则运行结果、指标、隔离样本"]

关键边界:

  • DataOps 是规则定义、版本、发布、绑定、权限、执行计划和审计的源真相。
  • Kestra 不保存业务规则正文或数据源凭据,只接收不可变版本 ID、计划 Hash、参数 和短期任务令牌。
  • Runner 不接受前端传来的任意 SQL/Python;只加载已发布的 ExecutionPlan。
  • AI 是自然语言规则的默认生产入口,但所有模型输出都先进入确定性控制器。控制器 决定校验、测试、发布、绑定和执行,模型不能绕过权限或直接连接生产数据源。
  • 模型不在运行时逐行判断数据;发布后的表达式、SQL、Polars 计划或签名代码由 Runner 确定性执行,因此同一版本可重放、可对账、可回滚。

AI Rule Authoring Agent 的闭环为:

  1. 保存用户原始描述,解析规则类型、输入输出、字段约束、异常处置、优先级和示例;
  2. 从元数据服务读取 Schema、样例的脱敏画像、目标方言和相似规则,禁止把凭据或 原始敏感数据发送给模型;
  3. 以封闭 JSON Schema 输出 RuleCandidate,包含 RuleSpec 或 GeneratedCodeSpec、假设、歧义、置信度、测试样例和可解释摘要;
  4. 控制器做字段/类型/函数/权限/影响范围校验,并自动生成单元、属性、golden 和 脱敏样本测试;失败时把结构化错误返回模型,最多自动修复两轮;
  5. 策略引擎决定自动发布,或因歧义、低置信度、破坏性写入、Schema 变更、全量回填 等原因请求澄清/审批;
  6. 记录模型提供方、模型名、Prompt 版本、Schema 版本、上下文 Hash、候选 Hash、 测试结果、修复轮次和最终决策,保证每次解释和生成可追溯。

两个设计入口复用上述闭环,但输出对象不同:

  • 数据标准入口:AI 先生成结构化 StandardCandidate(适用对象、标准条款、严重度、 生效范围和例外),再为每条可执行条款生成或复用 RuleCandidate。发布 StandardVersion 时固定所有 RuleVersion,并能回答“哪些生产线正在执行/尚未执行这条标准”。
  • 数据流程入口:AI 生成 DataFlowCandidate(阶段、输入输出、标准组件、规则组件、 顺序和质量闸门)。它优先搜索并装配已有 StandardVersion/RuleVersion;确实没有可复用 资产时才创建 flow_scoped 规则并走同一生成闭环,避免流程内藏匿名脚本。
  • 数据工厂入口:不生成或修改业务语义,只选择已发布生产线版本并配置环境、调度、 资源和上线策略。任何规则变化都先回到标准/流程设计态生成新版本,再重新投产。

4.3 RuleSpec 设计

RuleSpec 采用 JSON,并通过 JSON Schema 严格关闭未知字段。示例:

{
  "schema_version": "1.0",
  "rule_uid": "019f0000-0000-7000-8000-000000000001",
  "name": "normalize_customer",
  "input_schema_ref": "bd:customer:v7",
  "output_schema_ref": "bd:customer_clean:v3",
  "steps": [
    {"id": "cast_age", "op": "cast", "column": "age", "to": "int64", "on_error": "quarantine"},
    {"id": "trim_name", "op": "normalize_text", "column": "name", "trim": true},
    {"id": "adult", "op": "derive", "target": "is_adult", "expression": "age >= 18"},
    {"id": "valid_id", "op": "assert", "expression": "customer_id != ''", "severity": "error", "on_failure": "reject"},
    {"id": "dedup", "op": "deduplicate", "keys": ["customer_id"], "keep": "last", "order_by": ["updated_at"]}
  ],
  "null_policy": "explicit",
  "timezone": "Asia/Shanghai"
}

首期允许的算子:

  • 类型与文本:castnormalize_textregex_replacefill_null
  • 行级逻辑:filterderivemap_valuesassert
  • 集合逻辑:deduplicateaggregatelookup_join
  • 敏感数据:mask(只允许平台预注册策略);
  • 结果路由:rejectquarantine,由 on_error/on_failure 引用。

所有表达式函数均使用 allowlist;时间、时区、空值、舍入、字符串大小写和正则语义 必须写入跨后端一致性测试。规则中禁止凭据、连接串、文件路径、网络地址、动态导入、 任意代码和多语句 SQL。

当 RuleSpec 算子不足以表达用户意图时,AI 可以生成 GeneratedCodeSpec,但不能把 聊天输出直接当脚本运行。该规范必须声明入口函数、输入输出 Schema、依赖 allowlist、 资源上限、副作用、测试和代码 Hash。构建服务对 Python AST、导入、文件/网络访问、 子进程和危险调用做静态检查,在无凭据沙箱中运行测试,生成包含依赖清单的不可变制品 并签名;Runner 只执行固定摘要的制品。优先扩展 RuleSpec 算子,代码生成是受控的 逃生舱,而不是默认捷径。

4.4 执行策略

编译器在发布时生成不可变 ExecutionPlan,而不是每次运行临时决定后端:

条件 执行后端 数据路径
同一关系库、算子全部可下推 sql_pushdown SQLGlot 生成严格方言 AST;在源/目标库内执行
跨数据源、文件或下推不支持 polars_batch 服务端游标分批读取,Parquet/Arrow 批次转换,幂等写入
仅质量验证 quality_check 内建断言;复杂 Suite 可委托 GX Adapter
RuleSpec 无法表达且代码策略允许 generated_python 构建并执行签名、固定摘要的受限 Python 制品
超过单机阈值或持续流 external_adapter 后续 Beam/Flink 适配器,本期明确拒绝而非静默降级

ExecutionPlan 至少保存:

  • rule_version_ids、规范 Hash、编译器版本和目标方言;
  • 输入/输出 Schema Hash、字段读写集合和血缘;
  • 后端、批大小、资源上限、流式兼容性和预估影响行数;
  • 读写数据源 UID、事务策略、幂等键、水位线和隔离目标;
  • 编译产物 Hash;SQL 只保存规范化 AST/模板,不保存凭据;生成代码保存签名制品引用、 镜像/运行时摘要和依赖清单。

4.5 数据生产线装配及与 WorkflowSpec / Runner 的集成

DataFlow 编排人员从目录中选择已发布的数据标准和数据规则(也可用自然语言让 AI 推荐 和组装),按阶段形成生产线,不再填写表达式、SQL 或脚本路径。发布 DataFlowVersion 时,Production Line Resolver 自动:

  1. 固定所选 standard_version_id,展开其 StandardRuleBinding 并固定全部 rule_version_id,禁止使用会漂移的 latest
  2. 合并直接选择的规则和标准展开规则,检测重复、冲突、顺序、字段读写和 Schema 契约;
  3. 为每个阶段选择 SQL pushdown、Polars、质量检查或签名代码后端,生成 ExecutionPlan;
  4. 生成不可变 ProductionLinePackage,包含 DataFlowVersion、组件快照、规则/标准版本、 ExecutionPlan Hash、输入输出 Schema、质量闸门和血缘;
  5. 将生产线标为 released,交给数据工厂投产;设计页面不得直接激活生产调度。

数据工厂选择这个 released DataFlowVersion 后,创建 DataFlowDeployment,绑定环境数据源、 秘密引用、资源限额、并发、调度、回填和告警策略,再自动生成 WorkflowSpec,以 disabled 状态部署 Kestra Flow,完成 dry-run/canary 后按策略激活。这样“生产线设计”和“生产线 投产”分别对应 DataFlow 与 Data Factory,规则语义不会在投产时被重新生成。

DataFlowSpec 中新增一个逻辑标准组件和两个物理规则节点类型:

{
  "dataflow_version_id": "019f0000-0000-7000-8000-000000000020",
  "components": [
    {
      "id": "customer_standard_gate",
      "type": "standard.enforce",
      "standard_version_id": "019f0000-0000-7000-8000-000000000030",
      "stage": "quality_gate"
    },
    {
      "id": "clean_customer",
      "type": "rule.apply",
      "rule_version_id": "019f0000-0000-7000-8000-000000000040",
      "stage": "transform",
      "idempotency": {
        "strategy": "partition_replace",
        "key": "customer:${parameters.biz_date}"
      }
    }
  ]
}
  • standard.enforce 是设计态逻辑节点;发布生产线时展开为固定版本的 quality.check/rule.apply 节点,但在部署包和运行结果中保留标准条款来源。
  • rule.apply 由 Runner 根据 dataflow_component_binding_id + plan_hash 读取已发布 计划,拒绝 Hash 不一致、未投产或已撤销版本。
  • 同库下推可在一个 rule.apply 中完成读取、变换和写入;跨库链路用 MinIO Parquet artifact_ref 传递批次,不能把数据行内联进 Kestra JSON。
  • quality.check 返回总行数、通过/失败数、失败比例、规则命中数和脱敏后的有限样本; 阈值决定 pass/warn/fail,再由 WorkflowSpec 决定是否继续。
  • Runner 的任务账本继续记录单次令牌和提交结果;规则运行表补充业务级指标与 rule_version_id
  • 预注册 Handler 与 AI 生成代码均必须满足同一制品策略;生成代码只有在自动静态 检查、沙箱测试、风险策略和签名全部通过后才可执行,高风险能力仍要求代码评审。

4.6 数据模型

建议在 PostgreSQL 新增以下逻辑表;可按迁移批次实施,但不能省略标准、生产线或投产 三类关系:

关键字段 说明
data_rules uid, name, category, owner_uid, status 稳定规则身份
data_rule_versions rule_uid, version_no, source_text, rule_spec, spec_hash, status, created_by 原始意图与不可变执行规范;draft/tested/published/deprecated
rule_generation_runs rule_version_id, model_provider/name, prompt/schema_version, context/candidate_hash, confidence, ambiguities, decision AI 解析、修复与策略决策审计
data_standards / data_standard_versions standard_uid, version_no, source_text, clauses, scope, status 数据标准稳定身份与不可变条款版本
standard_rule_bindings standard_version_id, clause_id, rule_version_id, severity, exception_policy 标准条款与可执行规则版本的固定关系
dataflow_versions dataflow_uid, version_no, input/output_schema_hash, status, package_hash 不可变数据生产线版本
dataflow_component_bindings dataflow_version_id, component_kind, standard/rule_version_id, stage, order_no 装配标准和规则,不保存复制的代码
rule_execution_plans component_binding_id, backend, compiler_version, plan, plan_hash, schema_hashes 生产线组件的发布期编译产物
rule_artifacts rule_version_id, kind, uri, digest, signature, runtime_digest, dependency_manifest 表达式、SQL、Polars 或生成代码制品
dataflow_deployments(由现有 dataflow_workflow_versions 演进) dataflow_version_id, environment, workflow_version_id, schedule_plan_id, status, activated_at/by 数据工厂投产、激活和回滚边界
rule_runs workflow_run_id, binding_id, rows_in/out/rejected/quarantined, status, timings 可审计执行结果
rule_violation_samples rule_run_id, artifact_ref, sample_count, redaction_policy 脱敏、限量的失败样本引用

Neo4j 只保留 Standard、Rule、DataFlow、BusinessDomain、DataProduct 的稳定 UID 及语义、 影响和血缘关系;不可变版本、组件顺序、发布/投产状态和运行结果放 PostgreSQL,避免 把事务状态拆散到图数据库。

4.7 API 与三个产品页面

建议新增后端接口:

  • 数据标准:POST /api/data-standards/interpretPOST .../{uid}/versionsPOST .../{version}/generate-rules|validate|simulate|publish
  • POST /api/data-rulesGET /api/data-rules
  • POST /api/data-rules/interpret:自然语言 + 上下文生成候选规则和歧义说明;
  • POST /api/data-rules/{uid}/versions
  • POST /api/data-rules/{uid}/versions/{version}/generate:生成/修复 RuleSpec 或代码;
  • POST /api/data-rules/{uid}/versions/{version}/validate
  • POST /api/data-rules/{uid}/versions/{version}/simulate:编译并在脱敏样本上试运行;
  • POST /api/data-rules/{uid}/versions/{version}/publish
  • 数据流程:POST /api/dataflows/{uid}/versionsPOST .../{version}/componentsPOST .../{version}/resolve|validate|simulate|release
  • 数据工厂:POST /api/data-factory/deploymentsPOST .../{id}/dry-run|canary|activate|rollbackGET .../{id}/runs
  • GET /api/rule-runs/{run_uid}

后端分别强制 standards:*rules:*dataflows:*data_factory:deploy/operate 权限。规则/标准发布、生产线 release 和生产激活是三个 独立审计动作;生产环境可要求职责分离,不能因为 AI 自动化而合并越权。

三个页面都使用 AI,但交付的功能不同:

  • 数据标准页面:以自然语言录入标准,展示条款拆解、适用范围、生成/复用的规则、 表达式或代码、样本合规结果,以及引用/未覆盖该标准的数据生产线。现有“操作代码” 文本框改为只读的“可执行实现”视图,不能把模型返回代码直接保存为标准。
  • 数据流程页面:以可视化生产线为主,支持从目录拖入 StandardVersion 和 RuleVersion, 或让 AI 根据流程描述推荐和组装;显示阶段、先后依赖、字段映射、标准覆盖率、规则冲突、 输入输出 Schema、模拟结果和 ProductionLinePackage。发布动作是“发布生产线版本”。
  • 数据工厂页面:从 released 生产线目录选择版本,配置环境数据源、调度、资源、并发、 告警和上线策略,查看 dry-run/canary、激活、回滚、执行批次、数据质量与产出数据产品; 不提供修改标准正文、规则正文或生成代码的入口。

完整用户旅程应是:在数据标准中输入“客户手机号去除空格后必须为 11 位数字”,AI 生成 标准条款和可执行规则;在数据流程中把该标准与“客户去重”“地区编码转换”规则按阶段 组装为“客户主数据生产线”;发布 DataFlowVersion 后,到数据工厂选择测试/生产环境并 投产。全程平台自动生成并固定执行实现,用户不复制代码,且每层都能看到自己的版本、 验证结果和审计记录。

5. 安全、正确性与运行约束

  1. 编译期失败关闭:未知算子、未知字段、方言不兼容、Schema 不匹配、秘密字段、 未注册 Handler 或非流式大任务全部拒绝发布。
  2. 数据库安全:SQL 使用 AST 和绑定参数;表/列来自已登记元数据;禁止自由表名、 多语句、DDL、存储过程和客户端自定义驱动。
  3. 写入安全:写规则必须声明 upsertpartition_replacededuplication_key;提交结果未知时禁止自动重放。
  4. 批次安全:服务端游标、批大小、行数/字节数、CPU/内存、超时和并发都有硬上限; 中间 Parquet 设置 TTL、加密、租户/环境前缀和生命周期清理。
  5. 语义一致性:每个算子建立 SQL/PostgreSQL、SQL/MySQL、Polars 的 golden conformance 数据集,覆盖 NULL、空串、Unicode、时区、DST、Decimal 和异常值。
  6. 隐私:运行日志不记录整行数据;失败样本默认脱敏、限量并走授权下载。
  7. 可观测性:记录规则命中率、拒绝率、隔离率、Schema 漂移、耗时、吞吐、重试、 水位线和编译计划版本,关联 correlation_id 与 Workflow Run。
  8. 模型输入安全:元数据、历史规则、日志和样例都视为不可信上下文,进行提示注入 隔离和脱敏;模型无数据源凭据、生产网络和发布权限,工具调用由控制器逐项授权。
  9. 生成代码安全:依赖 allowlist、AST/污点扫描、无凭据沙箱、CPU/内存/时限、 只读文件系统、默认断网、制品签名和摘要校验缺一不可;不得在 Runner 动态安装依赖。

自动化策略建议按环境和影响分级:

等级 典型规则 默认动作
低风险 只读校验、格式标准化、测试环境派生字段 校验和测试通过后自动发布、绑定、canary
中风险 有幂等键的清洗写入、小范围分区替换 通过影响阈值与 canary 后可按环境策略自动提升
高风险 DDL、Schema 变更、全量覆盖/回填、不可逆删除、敏感数据外发 强制人工审批;部分能力直接禁止
语义不确定 多字段可能匹配、约束冲突、置信度低、缺少异常策略 请求用户澄清,不自动猜测和执行

6. 迁移策略

现有链路不能一次性切断:

  1. 同时盘点 data_standard.describe/codeDataFlow.script_requirement.rule、现有 Python 脚本和数据工厂工作流;分别保留为 legacy_source_text/code/workflow_ref 以便回滚。
  2. 批量 AI 分析标准描述、流程规则和已有脚本,生成 StandardCandidate/RuleCandidate, 保留模型/Prompt/上下文 Hash、假设、歧义和测试;不把标准代码与流程代码分成两套。
  3. 先发布 DataStandardVersion、RuleVersion 和 StandardRuleBinding;对标准页面现有 “操作代码”改为新制品的只读引用,禁止继续产生未版本化代码。
  4. 将每个活跃 DataFlow 转为 DataFlowVersion,装配已固定的标准和规则,生成 ProductionLinePackage;对缺少标准覆盖的关键输出给出告警而不是静默遗漏。
  5. 在数据工厂创建测试环境 DataFlowDeployment,以 shadow/canary 对账行数、字段 Hash、 聚合值、质量规则和拒绝样本;达到连续窗口门槛后再激活生产部署。
  6. 所有活跃标准、生产线和部署完成迁移且观察期通过后,才删除“模型文本代码”和 “自由文本 -> 任意脚本 -> SSH”的运行依赖;历史定义、部署和审计继续保留。

7. 分阶段实施计划

R0:契约与 PoC

  • 定义 StandardCandidate、RuleCandidate、DataFlowCandidate、RuleSpec 1.0、 GeneratedCodeSpec、ProductionLinePackage、算子语义、NULL/时区策略和封闭 JSON Schema。
  • 基于现有 orchestration/agent 实现 Rule Authoring Agent PoC,验证结构化生成、 错误反馈、有界修复、模型/Prompt/Hash 审计和不确定性输出。
  • 完成 CEL Python、SQLGlot PostgreSQL/MySQL 和 Polars Streaming 三个 PoC。
  • 建立同一规则跨后端 golden conformance 测试。
  • 验收:自然语言样例可稳定生成相同语义的结构化候选;相同输入在 PostgreSQL、MySQL、 Polars 得到一致结果;歧义和不支持语义明确失败。

R1:AI 数据标准与规则控制面

  • 建立 DataStandardVersion、RuleVersion、StandardRuleBinding、生成审计、Repository、 RBAC,以及解释/生成/校验/试运行/发布 API。
  • 数据标准页交付自然语言主入口、条款拆解、AI 理解与歧义、生成/复用规则、表达式/代码 只读制品视图、版本 Diff 和样本预览;结构化表单与 JSON 作为高级模式。
  • 验收:标准条款全部可追到规则版本;低风险规则可自动发布;歧义请求澄清;已发布版本 不可修改;旧页面不能直接保存未经治理的模型代码。

R2:数据流程生产线装配

  • 建立 DataFlowVersion、standard.enforce/rule.apply/quality.check 组件和 Production Line Resolver,自动展开标准、固定规则、检测冲突并生成部署包。
  • 数据流程页支持标准/规则目录、可视化阶段装配、AI 推荐、Schema 契约、覆盖率、模拟、 版本 Diff 和 release;此阶段不直接激活 Kestra Flow。
  • 验收:一个生产线同时装配标准和规则;发布后所有引用均为不可变版本;标准覆盖缺失、 规则冲突和字段契约失败会阻止 release;不需要手写脚本或调度节点。

R3:数据工厂投产与执行面

  • 建立 DataFlowDeployment;数据工厂实现环境、数据源、调度、资源、dry-run、canary、 activate、rollback 和运行监控,确定性生成 WorkflowSpec/Kestra Flow。
  • Runner 增加 SQL pushdown、Polars batch、MinIO artifact 和隔离输出执行器;把 sql.query 改为服务端游标/分页,禁止 .all() 全量载入。
  • 上线 GeneratedCodeSpec 构建沙箱、静态检查、自动测试、签名制品和 Runner 摘要校验; 在此之前,生产环境只允许 RuleSpec 编译产物和预注册 Handler。
  • 评估并按需接入 Great Expectations Adapter。
  • 验收:PostgreSQL/MySQL 跨库 10 万行可分批处理且幂等;每次运行可追到 Deployment、 DataFlowVersion、StandardVersion、RuleVersion、Plan Hash、Schema 和数据产品;失败可 安全重放或明确标为 unknown

R4:历史迁移与退出

  • 迁移活跃 data_standard.codescript_requirement 和数据工厂工作流,逐条对账、 切换和回滚演练。
  • 接入规则/标准运行指标、失败样本脱敏、Schema 漂移、标准覆盖和数据产品合格证。
  • 清理硬编码 SSH Credential ID 和任意脚本运行入口。
  • 验收:连续观察窗口无差异、无活跃标准/生产线/部署依赖旧代码或 SSH 路径,审计与 回滚证据齐全。

后续可选

  • 数据量或实时需求超过单机阈值时增加 Beam/Flink 执行 Adapter;不改变 RuleSpec。
  • 高价值复杂插件可评估签名容器或 WebAssembly 沙箱;仍不允许页面上传任意代码。
  • 只有出现复杂事实推理、DMN 决策表或 CEP 需求时才引入 Drools Sidecar。

8. 建议的首个垂直切片

选择“客户主数据”场景,跨越三个产品页面:

  1. 数据标准输入“手机号去空格后必须为 11 位数字;不合格数据隔离”,AI 生成 StandardVersion、RuleVersion、测试和 StandardRuleBinding;
  2. 数据流程把该标准与“类型转换 + 地区编码转换 + 客户去重”规则按 normalize -> transform -> quality_gate -> write 装配成 DataFlowVersion;
  3. Resolver 展开标准、固定全部规则版本、生成 Polars/SQL ExecutionPlan 和 ProductionLinePackage,模拟通过后 release,但不在此处启动生产;
  4. 数据工厂选择该生产线版本,绑定测试数据源、每天调度和资源上限,生成 disabled Kestra Flow,执行 dry-run/canary 后激活;
  5. Runner 产生 pass/reject/quarantine、标准符合率、脱敏样本、血缘和数据产品合格证, 重复触发验证幂等;规则升级后验证旧部署不漂移,新部署可独立回滚;
  6. 用一条故意含糊的标准描述验证平台请求澄清,并验证数据流程不能 release 缺少必要 StandardVersion 的生产线。

该切片验证“数据标准定义要求 -> AI 生成执行实现 -> 数据流程装配生产线 -> 数据工厂 上线投产 -> 生产数据产品”的完整闭环,能直接防止只实现规则页而遗漏标准绑定、生产线 装配或工厂投产能力。

9. 参考资料