ADR-004-ai-first-kestra-orchestration.md 5.9 KB

ADR-004:AI-first Data Factory 与 Kestra 调度边界

背景

DataOps Platform 当前通过 n8n 定义、调度和执行数据加工流程。n8n 使用带有 使用限制的 Sustainable Use License,且当前集成把平台业务模型绑定到了 n8n_workflow_id、n8n Workflow JSON 和 SSH 脚本执行方式。

平台已经具备稳定 DataFlow UID、环境级版本激活、RBAC、Outbox、外部数据源 加密凭据和 Worker 本地连接池。后续目标不是继续由人员在第三方调度 UI 中维护 流程,而是由调度智能体依据血缘、SLA、运行历史和资源状态进行规划,通过 MCP 完成受控发布、执行、恢复和优化。

决策

  1. DataOps 是工作流定义、版本、权限、发布状态、业务上下文和审计记录的源真相。
  2. Kestra OSS 是首选调度执行引擎,负责确定性调度、任务状态机、重试、回填、 重放和执行日志;Kestra 不是 DataFlow 治理对象的源真相。
  3. 调度智能体不直接连接拥有完整权限的 Kestra MCP。DataOps 建设 Scheduling MCP Gateway,只暴露经过权限和策略约束的复合操作。
  4. DataOps 建设 Context MCP,向智能体提供 DataFlow、血缘、SLA、数据源能力、 连接池健康和历史运行指标,但不返回用户名、密码、密文或连接串。
  5. 智能体生成平台自有的 WorkflowSpecSchedulePlan;确定性编译器将通过 校验的版本转换为 Kestra Flow。不得把未经校验的模型输出直接发布到生产。
  6. 数据任务由独立 DataOps Runner 执行。Kestra 只传递 DataFlow/任务版本、 data_source_uid、用途和参数;Runner 复用现有 DataSourceConnectionManager,Kestra 不保存业务数据源凭据。
  7. 日常规划和恢复允许无人值守,但必须经过机器可验证的权限、风险、连接预算、 幂等、时间窗口和回填范围策略。高风险动作未满足策略时自动拒绝,不由模型绕过。
  8. n8n 与 Kestra 在迁移期并行运行。按单个 DataFlow 完成影子运行、结果对账、 切换和可回滚验证;达到统一退出门槛后才下线 n8n。
  9. 工作流引擎使用可插拔适配器接口。DataOps 持久化 engine_type 和外部绑定, 不再把任一引擎名称写入新的领域字段名。

目标架构

flowchart LR
    GOAL["业务目标 / SLA / 事件"] --> AGENT["调度规划智能体"]
    AGENT --> CTX["DataOps Context MCP"]
    AGENT --> GW["Scheduling MCP Gateway"]
    CTX --> GRAPH["DataFlow / 血缘 / 运行历史"]
    CTX --> HEALTH["数据源与连接池健康"]
    GW --> POLICY["确定性策略与权限校验"]
    POLICY --> SPEC["WorkflowSpec / SchedulePlan"]
    SPEC --> COMPILER["Kestra 编译器与引擎适配器"]
    COMPILER --> KMCP["受限 Kestra MCP"]
    KMCP --> KESTRA["Kestra OSS"]
    KESTRA --> RUNNER["DataOps Runner"]
    RUNNER --> POOL["DataSourceConnectionManager"]
    POOL --> SOURCE["PostgreSQL / MySQL 数据源"]
    KESTRA --> EVENTS["执行事件、日志和指标"]
    EVENTS --> AGENT
    EVENTS --> AUDIT["DataOps 运行审计"]

领域边界

能力 源真相/责任方
DataFlow 业务定义、输入输出和血缘 Neo4j + DataOps API
WorkflowSpec、版本、环境激活、策略和审计 DataOps PostgreSQL
Cron/Event 调度、执行状态机、重试和回填 Kestra OSS
AI 目标分解、计划生成和运行优化 DataOps 调度智能体
MCP 工具授权、参数约束和复合事务 DataOps Scheduling MCP Gateway
数据库凭据、连接池、事务和熔断 DataOps Runner
大文件、中间产物和归档日志 MinIO

AI 与确定性执行边界

智能体可以:

  • 读取依赖、SLA、历史耗时、失败原因和资源状态。
  • 生成或调整候选流程、计划时间、并发度、重试和回填建议。
  • 发布禁用版本、执行 Canary、在策略允许时推广或回滚。
  • 对已知错误执行限次重试、降并发、暂停或切换前一稳定版本。

智能体不可以:

  • 读取或修改数据源明文凭据。
  • 绕过 WorkflowSpec、JSON Schema、RBAC 或策略校验直接发布 Kestra YAML。
  • 无边界删除流程/执行记录、修改任意状态或发起无限区间回填。
  • 自动重放没有幂等证明的数据写入任务。
  • 把日志、元数据描述或外部数据中的文本当作可信系统指令。

Kestra 开源边界

本方案只依赖 Kestra Apache 2.0 开源核心和 Apache 2.0 官方 Python MCP Server。 Kestra 企业版的 RBAC、服务账号、审计和秘密管理不作为平台依赖;这些能力由 DataOps 已有身份权限和新建 MCP Gateway 提供。Kestra 只允许在内部网络被 Gateway 和平台适配器访问。

迁移兼容

  • 现有 dataflow_workflow_versions 保留,先添加通用引擎字段并回填 engine_type='n8n',不立即删除 n8n_workflow_id
  • 迁移期间 N8nAdapterKestraAdapter 同时存在。
  • 同一个 DataFlow/环境只有一个正式执行引擎;影子引擎只能写入隔离目标或执行 只读校验,避免重复生产写入。
  • n8n 下线后保留定义归档、映射和必要执行历史,最后再通过独立迁移移除运行依赖。

结果

  • DataOps 获得面向智能体的调度控制面,同时保持底层执行确定性。
  • Kestra 可以被替换而不改变 DataFlow 领域模型或数据库资源管理。
  • 日常调度、故障恢复和优化可以无人值守,但所有变更可验证、可审计、可回滚。
  • 增加 DataOps MCP、Runner、策略引擎和双轨迁移的建设成本。

详细实施工作包和退出门槛见 Kestra AI-first Data Factory 双轨改造实施计划