# ADR-004:AI-first Data Factory 与 Kestra 调度边界 - 状态:Accepted - 日期:2026-07-18 - 替代:[ADR-002:DataFlow 与 n8n Workflow 的职责和版本映射](ADR-002-workflow-engine.md) ## 背景 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. 智能体生成平台自有的 `WorkflowSpec` 和 `SchedulePlan`;确定性编译器将通过 校验的版本转换为 Kestra Flow。不得把未经校验的模型输出直接发布到生产。 6. 数据任务由独立 `DataOps Runner` 执行。Kestra 只传递 DataFlow/任务版本、 `data_source_uid`、用途和参数;Runner 复用现有 `DataSourceConnectionManager`,Kestra 不保存业务数据源凭据。 7. 日常规划和恢复允许无人值守,但必须经过机器可验证的权限、风险、连接预算、 幂等、时间窗口和回填范围策略。高风险动作未满足策略时自动拒绝,不由模型绕过。 8. n8n 与 Kestra 在迁移期并行运行。按单个 DataFlow 完成影子运行、结果对账、 切换和可回滚验证;达到统一退出门槛后才下线 n8n。 9. 工作流引擎使用可插拔适配器接口。DataOps 持久化 `engine_type` 和外部绑定, 不再把任一引擎名称写入新的领域字段名。 ## 目标架构 ```mermaid 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`。 - 迁移期间 `N8nAdapter` 和 `KestraAdapter` 同时存在。 - 同一个 DataFlow/环境只有一个正式执行引擎;影子引擎只能写入隔离目标或执行 只读校验,避免重复生产写入。 - n8n 下线后保留定义归档、映射和必要执行历史,最后再通过独立迁移移除运行依赖。 ## 结果 - DataOps 获得面向智能体的调度控制面,同时保持底层执行确定性。 - Kestra 可以被替换而不改变 DataFlow 领域模型或数据库资源管理。 - 日常调度、故障恢复和优化可以无人值守,但所有变更可验证、可审计、可回滚。 - 增加 DataOps MCP、Runner、策略引擎和双轨迁移的建设成本。 详细实施工作包和退出门槛见 [Kestra AI-first Data Factory 双轨改造实施计划](../superpowers/plans/2026-07-18-kestra-ai-first-data-factory-migration.md)。