# Kestra AI-first Data Factory 双轨改造实施计划 > 状态:规划已确认,实施顺序固定为 V50 → V55。执行时每个 Task 均采用测试 > 先行,并在进入下一版本前提交对应的可重复验证证据。版本顺序和增量测试范围见 > [V50–V55 顺序交付与增量验证计划](2026-07-18-kestra-v50-v55-delivery-plan.md)。 **目标:** 在不中断现有 n8n 生产流程的前提下,建立 DataOps 自有的 AI 调度 控制面、MCP 安全网关和数据任务运行时,以 Kestra OSS 作为确定性调度引擎; 完成逐流程双跑、对账、切换和回滚验证后下线 n8n。 **架构:** DataOps 保存 WorkflowSpec、版本、策略、环境激活和审计;调度智能体 通过 Context MCP 获取业务上下文,通过 Scheduling MCP Gateway 执行受控复合 操作;Kestra 负责调度状态机,DataOps Runner 使用现有数据库资源池执行数据任务。 **技术栈:** Flask、PostgreSQL 16、Neo4j、MinIO、Kestra OSS、Kestra Python MCP Server、DeepSeek/可配置 LLM、SQLAlchemy、Docker Compose、pytest、Vue 2。 --- ## 1. 成功标准 ### 1.1 功能标准 - 智能体能从业务目标、SLA、DataFlow 血缘和数据源健康生成结构化调度计划。 - 支持手动、Cron、指定时间和事件触发。 - 支持依赖、条件、并行、子流程、超时、限次重试、暂停、终止、恢复、Replay、 Backfill 和版本回滚。 - 首批节点支持 `sql.query`、`sql.execute`、`python`、`http`、`condition`、 `parallel`、`subflow` 和 `notify`。 - 数据库节点只引用 `data_source_uid`,不在 Kestra 或 WorkflowSpec 中保存凭据。 - 开发、测试、生产环境仍保持一个 DataFlow 只有一个当前正式生效版本。 - n8n 和 Kestra 可以按 DataFlow 独立选择正式/影子角色,不进行全局一次性切换。 ### 1.2 AI 自治标准 - 常规计划生成、测试、推广、监控和可恢复故障处理不依赖人员操作。 - 模型输出必须通过 JSON Schema、DAG、权限、连接预算、幂等和风险策略校验。 - 相同输入和固定策略下,执行计划的可变部分有明确边界,调度执行本身保持确定性。 - 模型不可用时,已发布的 Kestra 调度继续运行;平台退化为固定计划,不中断生产。 - 所有 AI 决策保存目标、上下文摘要、模型/提示版本、候选计划、校验结果和动作结果。 ### 1.3 运维标准 - DataOps、Kestra、Runner、MCP Gateway 任一组件重启不产生重复生产写入。 - 单个数据源故障不影响其他数据源、平台控制库或无关流程。 - 关键动作具有关联 ID,可从 DataFlow 追踪到计划、Kestra 执行、Runner 任务和 数据源池指标。 - n8n 切换失败时可在约定恢复时间内回到上一正式版本。 ## 2. 并行工作流与依赖 ```mermaid flowchart LR K0["K0 决策冻结与盘点"] --> A["A 引擎无关领域模型"] K0 --> B["B Kestra 与引擎适配"] K0 --> C["C DataOps Runner"] K0 --> D["D Context MCP"] A --> E["E Scheduling MCP Gateway"] B --> E D --> E C --> F["F 首批数据任务节点"] E --> G["G 调度规划智能体"] F --> G G --> H["H 双轨影子运行与对账"] H --> I["I 分批切换与稳定观察"] I --> J["J n8n 只读归档与下线"] ``` 可以并行推进: - A、B、C、D 在 K0 完成后并行。 - Kestra 本地环境、Runner 解耦和 Context MCP 契约可由不同开发线独立推进。 - 流程迁移器和结果对账框架可在首批节点实现后并行扩展。 必须串行的门槛: - 未完成策略校验和禁用版本发布前,不允许智能体操作生产 Kestra。 - 未完成 Runner 幂等/事务验证前,不允许影子流程写入正式目标。 - 未完成双跑对账和回滚演练前,不允许切换正式引擎。 - 未满足统一退出门槛前,不允许删除 n8n 运行依赖或历史数据。 ## 3. Task 0:冻结决策、资产盘点与基线 **文件:** - 修改:`docs/architecture/ADR-002-workflow-engine.md` - 创建:`docs/architecture/ADR-004-ai-first-kestra-orchestration.md` - 创建:`docs/generated/n8n_migration_inventory.json` - 创建:`scripts/inventory_n8n_workflows.py` - 创建:`tests/test_n8n_migration_inventory.py` **工作项:** - [ ] 导出全部 n8n Workflow 的 ID、名称、启停状态、触发器、节点类型、凭据引用名、 Webhook、最近执行状态和关联 DataFlow;不导出凭据内容。 - [ ] 把流程分为 SQL、Python/SSH、HTTP、条件、通知、AI、未知/社区节点等类别。 - [ ] 标记生产写入、是否幂等、是否可影子执行、目标表/对象和业务 SLA。 - [ ] 为每个流程确定迁移优先级、负责人、兼容性状态和回滚路径。 - [ ] 固化 n8n 基线定义哈希,后续迁移期检测未经 DataOps 管理的外部修改。 **完成证据:** - 资产清单覆盖全部活跃和近 90 天有执行记录的流程。 - 清单不含 API Key、密码、Token、连接串或 n8n Credential 内容。 - 未知节点和无 DataFlow UID 的流程均进入阻塞清单,不被静默忽略。 ## 4. Task 1:建立引擎无关工作流领域模型 **文件:** - 创建:`migrations/versions/20260718_60_workflow_engine_abstraction.py` - 创建:`app/core/orchestration/models.py` - 创建:`app/core/orchestration/spec.py` - 创建:`app/core/orchestration/repository.py` - 创建:`app/core/orchestration/policy.py` - 修改:`app/core/data_flow/workflow_repository.py` - 创建:`tests/core/orchestration/test_spec.py` - 创建:`tests/core/orchestration/test_repository.py` - 创建:`tests/core/orchestration/test_policy.py` **工作项:** - [ ] 为现有版本记录增加 `engine_type`、`engine_definition_id`、 `engine_revision` 和 `deployment_metadata`,回填现有记录为 `n8n`。 - [ ] 迁移期保留 `n8n_workflow_id` 和 `n8n_workflow_name`,禁止破坏性删除。 - [ ] 建立 `workflow_schedules`、`workflow_runs`、`workflow_task_runs`、 `workflow_engine_bindings` 和 `workflow_plan_audits`。 - [ ] 定义版本化 WorkflowSpec JSON Schema,节点只允许注册表中的固定类型。 - [ ] 定义 SchedulePlan:时区、Cron/事件、SLA、最大并发、冲突策略、重试、 Backfill 上限和资源预算。 - [ ] 延续环境级唯一正式版本约束,并增加正式引擎角色与影子引擎角色约束。 - [ ] 保留 canonical hash 和秘密字段递归脱敏。 **完成证据:** - 现有 n8n 版本可无损读取和回填。 - 同一 DataFlow/环境不能同时存在两个正式引擎绑定。 - 非注册节点、环形 DAG、明文秘密、无限重试和无界回填均被拒绝。 ## 5. Task 2:部署本地 Kestra OSS 与官方 MCP **文件:** - 修改:`deploy/docker/docker-compose.yml` - 修改:`deploy/docker/postgres/init/000-init.sql` - 修改:`deploy/docker/README.md` - 创建:`deploy/docker/kestra/application.yml` - 创建:`deploy/docker/kestra/mcp.env.example` - 创建:`tests/test_kestra_local_contract.py` - 创建:`tests/integration/test_kestra_mcp_contract.py` **工作项:** - [ ] 使用固定版本镜像,Kestra 元数据使用独立数据库或独立 Schema。 - [ ] Kestra 和官方 Python MCP 仅加入内部网络,不直接暴露生产公网入口。 - [ ] OSS MCP 配置禁用 `ee` 工具组,并由 Gateway 控制允许的工具集合。 - [ ] 验证 Flow 创建/更新、启停、执行、日志、Pause/Kill、Backfill、Replay、 Restart 和 Resume 的实际契约。 - [ ] 为 MCP 版本升级建立工具清单和参数 Schema 差异测试。 - [ ] Kestra 不创建业务数据库凭据,不启用直接连接生产数据源的任务。 **完成证据:** - 本地 Compose 可重复启动,Kestra 和 MCP 健康检查通过。 - MCP 合约测试覆盖所有计划使用的写操作。 - Kestra 或 MCP 不可用时,DataOps 控制面和已存在的数据源管理仍可用。 ## 6. Task 3:实现可插拔引擎适配器和 Kestra 编译器 **文件:** - 创建:`app/core/orchestration/engines/base.py` - 创建:`app/core/orchestration/engines/n8n.py` - 创建:`app/core/orchestration/engines/kestra.py` - 创建:`app/core/orchestration/compilers/kestra.py` - 创建:`app/core/orchestration/service.py` - 创建:`tests/core/orchestration/test_engine_contract.py` - 创建:`tests/core/orchestration/test_kestra_compiler.py` **工作项:** - [ ] 定义 `validate`、`deploy_disabled`、`activate`、`deactivate`、`execute`、 `pause`、`kill`、`replay`、`backfill`、`get_execution` 和 `get_logs` 接口。 - [ ] 用现有 N8nClient 包装 `N8nAdapter`,保持迁移期行为兼容。 - [ ] 将 WorkflowSpec 确定性编译为 Kestra YAML,相同输入产生相同定义哈希。 - [ ] 所有 Flow 带 DataFlow UID、环境、版本、计划 ID 和 correlation ID 标签。 - [ ] 编译器默认生成禁用的 Schedule Trigger,只有推广事务可以启用。 - [ ] 在激活失败时写入可重试状态和 Outbox,不直接修改 Neo4j 血缘定义。 **完成证据:** - 引擎合约测试对 N8nAdapter 和 KestraAdapter 使用相同测试集。 - 编译快照测试证明同一 WorkflowSpec 的输出稳定。 - Kestra 外部成功、PostgreSQL 最终提交失败时可通过 reconcile 恢复。 ## 7. Task 4:把数据库资源池解耦为 DataOps Runner **文件:** - 创建:`app/runner/__init__.py` - 创建:`app/runner/api.py` - 创建:`app/runner/executor.py` - 创建:`app/runner/auth.py` - 创建:`app/runner/node_registry.py` - 修改:`app/core/data_source/runtime.py` - 修改:`app/core/data_source/manager.py` - 创建:`deploy/docker/runner.Dockerfile` - 修改:`deploy/docker/docker-compose.yml` - 创建:`tests/runner/test_executor.py` - 创建:`tests/integration/test_runner_datasource_pool.py` **工作项:** - [ ] 将连接池构造从 Flask `current_app` 解耦为显式配置,使 Runner 可独立启动。 - [ ] Kestra 只向 Runner 传递签名任务令牌、版本、节点 ID、`data_source_uid`、 `purpose`、参数和 correlation ID。 - [ ] 任务令牌短时有效、单任务绑定、不可重放,不包含数据源凭据。 - [ ] SQL 模板采用参数化语句;`sql.query` 强制只读,`sql.execute` 需要写入策略。 - [ ] 写节点记录幂等键、目标和提交结果;连接失败不自动重放未知状态的写事务。 - [ ] Python 节点使用固定镜像/依赖白名单、CPU/内存/时间和网络限制。 - [ ] Runner 副本数、Worker 数和池参数纳入全局连接预算。 **完成证据:** - Kestra 数据库中不存在业务数据源密码或完整连接串。 - Runner 重启、超时和重复请求不会导致已提交写任务重复执行。 - 单数据源熔断不影响其他数据源和平台控制库。 ## 8. Task 5:建设 DataOps Context MCP **文件:** - 创建:`mcp-servers/dataops-context/server.py` - 创建:`mcp-servers/dataops-context/tools/dataflows.py` - 创建:`mcp-servers/dataops-context/tools/lineage.py` - 创建:`mcp-servers/dataops-context/tools/datasources.py` - 创建:`mcp-servers/dataops-context/tools/executions.py` - 创建:`mcp-servers/dataops-context/tools/plans.py` - 创建:`tests/mcp/test_dataops_context_mcp.py` **只读工具:** - `list_dataflows` - `describe_dataflow` - `get_dataflow_dependencies` - `get_data_lineage` - `list_datasource_capabilities` - `get_datasource_pool_health` - `get_execution_history` - `get_sla_constraints` - `estimate_schedule_capacity` **校验工具:** - `validate_workflow_spec` - `validate_schedule_plan` - `simulate_schedule` - `compare_execution_results` **工作项:** - [ ] 工具返回结构化、限量和权限过滤后的数据,禁止返回秘密和大体量业务数据。 - [ ] 日志、描述、元数据值和外部错误统一标记为不可信内容,不能形成系统指令。 - [ ] 所有查询绑定智能体身份、角色、业务域和 correlation ID。 - [ ] 给每个工具定义明确的最大行数、时间范围、超时和错误结构。 **完成证据:** - Viewer/Editor/调度智能体只能看到授权业务域。 - 数据源工具只返回 UID、类型、用途、健康和容量摘要。 - Prompt injection 测试证明日志或元数据文本不能改变工具权限和策略。 ## 9. Task 6:建设 Scheduling MCP Gateway 和机器策略 **文件:** - 创建:`mcp-servers/dataops-scheduling/server.py` - 创建:`mcp-servers/dataops-scheduling/policy.py` - 创建:`mcp-servers/dataops-scheduling/kestra_client.py` - 创建:`mcp-servers/dataops-scheduling/tools/planning.py` - 创建:`mcp-servers/dataops-scheduling/tools/lifecycle.py` - 创建:`tests/mcp/test_scheduling_gateway.py` - 创建:`tests/security/test_scheduling_tool_boundaries.py` **向智能体开放的复合工具:** - `create_candidate_plan` - `deploy_disabled_version` - `run_canary` - `promote_candidate` - `pause_schedule` - `retry_failed_execution` - `backfill_bounded_window` - `rollback_to_previous_version` **不得直接开放:** - 原始 `delete flow/execution` - 任意 `change_status` - 任意 Namespace 文件和 KV 修改 - 动态 Kestra 连接地址/API Key 修改 - 未经校验的 YAML 直接创建生产 Flow - 无最大范围的 Backfill 或无次数限制的重试 **策略至少校验:** - 身份、环境和业务域权限。 - WorkflowSpec/Plan Schema、DAG 和节点注册表。 - 数据源用途、读写权限和连接预算。 - 写入幂等证明、目标隔离和重复执行风险。 - 最大并发、超时、重试、Backfill 窗口和成本上限。 - 生产变更窗口、Canary 结果、当前稳定版本和回滚目标。 **完成证据:** - 直接调用危险 Kestra MCP 工具的请求被网关拒绝并审计。 - 同一推广请求重复提交保持幂等。 - 策略拒绝原因可由智能体理解,但不泄漏凭据或内部安全细节。 ## 10. Task 7:实现调度规划智能体 **文件:** - 创建:`app/core/orchestration/agent/planner.py` - 创建:`app/core/orchestration/agent/prompts.py` - 创建:`app/core/orchestration/agent/schemas.py` - 创建:`app/core/orchestration/agent/evaluator.py` - 创建:`app/core/orchestration/agent/recovery.py` - 创建:`tests/agent/test_planner_scenarios.py` - 创建:`tests/agent/test_recovery_scenarios.py` **工作项:** - [ ] 输入为业务目标、SLA、可用时间窗和授权资源,输出严格结构化 Plan。 - [ ] 规划顺序固定为:观察、生成候选、校验、模拟、禁用发布、Canary、推广、 监控、恢复/回滚。 - [ ] 模型只决定允许的可变参数;Cron 触发和已发布运行不依赖模型在线。 - [ ] 保存模型、Prompt、Schema、上下文哈希和决策摘要,支持离线重放评估。 - [ ] 对模型超时、无效 JSON、工具失败和相互冲突的目标提供确定性降级。 - [ ] 恢复智能体只能执行策略允许的暂停、限次重试、降并发和回滚。 **固定评测场景:** - 上游延迟但必须在 SLA 前完成。 - 数据源池退化,需要降低并发或错峰。 - 同一目标表存在两个潜在写流程。 - 补跑窗口包含已经成功的分区。 - 模型建议无限重试、删除执行记录或读取密码。 - Kestra、Runner、模型或单个数据源分别不可用。 **完成证据:** - 固定场景评测全部产生合法计划或安全拒绝。 - 模型离线时已发布流程继续按确定性计划执行。 - 不同模型或模型版本不能绕过相同策略边界。 ## 11. Task 8:双轨运行、结果对账和单流程切换 **文件:** - 创建:`app/core/orchestration/migration/dual_run.py` - 创建:`app/core/orchestration/migration/reconciliation.py` - 创建:`app/commands/migrate_n8n_workflow.py` - 创建:`app/commands/reconcile_dual_run.py` - 创建:`tests/integration/test_n8n_kestra_dual_run.py` - 创建:`docs/runbooks/kestra-dual-run-and-rollback.md` **双轨模式:** | 模式 | 正式引擎 | 影子引擎 | 写入规则 | |---|---|---|---| | `n8n_primary` | n8n | 无 | 原有行为 | | `n8n_primary_kestra_shadow` | n8n | Kestra | 只读或隔离目标 | | `kestra_primary_n8n_standby` | Kestra | n8n | 仅 Kestra 写正式目标 | | `kestra_primary` | Kestra | 无 | n8n 只读归档 | **工作项:** - [ ] 将兼容的 n8n Workflow 转换为 WorkflowSpec,再编译为 Kestra Flow。 - [ ] 影子执行使用相同输入快照,但写入隔离 Schema/表/对象键。 - [ ] 对账行数、主键集合、聚合值、内容哈希、异常记录、耗时和资源消耗。 - [ ] 为允许的非确定性字段定义显式忽略/容差规则。 - [ ] 单流程切换使用事务化角色变更和 Outbox,失败自动恢复前一角色。 - [ ] 对无法自动迁移的社区节点生成阻塞报告,不自动替换为任意脚本。 **完成证据:** - 影子流程不会重复写正式目标。 - 对账报告可定位到具体节点、分区和差异类型。 - 每个切换流程均完成 Kestra 故障和回切 n8n 演练。 ## 12. Task 9:分批推广和稳定性观察 **推广批次:** 1. 只读、低频、无下游依赖流程。 2. 幂等写入、可按分区覆盖的批处理流程。 3. 多节点、有下游依赖的核心流程。 4. 高风险写入、长任务和特殊节点流程。 **每批进入条件:** - 资产清单、WorkflowSpec、策略和回滚目标完整。 - 至少完成约定次数的影子运行,覆盖正常和一次失败恢复。 - 对账结果达到该流程预先声明的相等/容差标准。 - 连接预算、峰值并发和 SLA 验证通过。 **每批退出条件:** - Kestra 正式运行持续达到约定观察窗口。 - 无未解释的数据差异、重复写入或关键 SLA 违约。 - 自动暂停、重试和回滚路径至少演练一次。 - n8n 已切换为 Standby 且未发生未授权外部修改。 ## 13. Task 10:n8n 统一退出门槛和下线 只有同时满足以下条件才允许进入下线变更: - [ ] 生产范围内所有 n8n Workflow 已归档、迁移或有正式保留说明。 - [ ] 所有正式 DataFlow 已使用 `kestra_primary`,不存在未处理双轨差异。 - [ ] 关键流程在完整业务周期内稳定运行,并覆盖一次恢复或回滚演练。 - [ ] n8n 已进入只读/Standby 观察窗口,期间无实际回切需求。 - [ ] Kestra、Runner、Gateway、Context MCP、智能体和对账指标均有监控与告警。 - [ ] 已验证模型不可用、Kestra 重启、Runner 重启和数据源故障不破坏生产状态。 - [ ] 已完成 n8n Workflow、执行历史、版本映射和凭据清单的安全归档。 - [ ] 已通过独立下线评审,确认没有 Webhook、前端入口、脚本或外部系统继续调用 n8n。 **下线文件:** - 修改:`deploy/docker/docker-compose.yml` - 修改:`deploy/docker/README.md` - 修改:`frontend/src/router/routes.js` - 修改:`frontend/src/api/dataFactory.js` - 修改:`app/api/data_factory/routes.py` - 修改:`app/core/data_factory/` - 创建:`migrations/versions/_retire_n8n_runtime.py` - 创建:`docs/runbooks/n8n-archive-and-retirement.md` **下线原则:** - 先停新建和激活,再停调度,再移除运行容器,最后清理代码和字段。 - 历史记录和归档先保留;`n8n_workflow_id` 的删除使用后续独立迁移。 - 禁止在同一个变更中同时删除 n8n、重构 WorkflowSpec 和升级 Kestra。 ## 14. 验证矩阵 | 维度 | 必验内容 | |---|---| | 合约 | WorkflowSpec、MCP 工具、Kestra API、Runner API、OpenAPI | | 数据 | 行数、主键、哈希、分区、重复写、事务不确定状态 | | 调度 | 时区、Cron、指定时间、并发、错过调度、Backfill、Replay | | 恢复 | Pause、Kill、Restart、Resume、回滚、组件重启 | | 安全 | RBAC、工具白名单、秘密脱敏、任务令牌、Prompt injection | | 资源 | 连接预算、池超时、熔断、Runner 扩缩、长任务 | | AI | 结构化输出、策略拒绝、模型降级、决策审计、场景回放 | | 迁移 | 双轨角色、影子隔离、对账、单流程切换、n8n 回切 | ## 15. V50–V55 交付节奏 本计划不以日历日期替代完成证据,按以下版本顺序组织: - **V50:** Task 0–1,建立迁移基线和引擎无关领域模型。 - **V51:** Task 2–3,部署 Kestra/MCP 并完成引擎适配和确定性编译。 - **V52:** Task 4,完成 DataOps Runner 和现有数据库资源池复用。 - **V53:** Task 5–6,完成 Context MCP、Scheduling MCP Gateway 和机器策略。 - **V54:** Task 7,打通受控 AI 规划闭环并执行首批只读 Canary。 - **V55:** Task 8–10,双轨对账、分批切换,并在满足门槛后有条件下线 n8n。 每个版本应独立形成可回滚变更;V50–V54 期间 n8n 均保持正式运行。V55 可以在 n8n Standby 状态结束阶段性交付,不把“已部署 Kestra”或“已经完成切换” 误报为“已经完成 n8n 下线”。 每个版本采用 L1 变更级、必要的 L2 边界级和 L3 场景级验证,不重复运行全量 测试。只有 V55 准备正式移除 n8n,或发生全局基础设施/公共协议变更时,才触发 一次 L4 全量验证。具体命令、退出条件和验证记录模板见 [V50–V55 顺序交付与增量验证计划](2026-07-18-kestra-v50-v55-delivery-plan.md)。