2026-07-18-kestra-ai-first-data-factory-migration.md 21 KB

Kestra AI-first Data Factory 双轨改造实施计划

状态:规划已确认,实施顺序固定为 V50 → V55。执行时每个 Task 均采用测试 先行,并在进入下一版本前提交对应的可重复验证证据。版本顺序和增量测试范围见 V50–V55 顺序交付与增量验证计划

目标: 在不中断现有 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.querysql.executepythonhttpconditionparallelsubflownotify
  • 数据库节点只引用 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. 并行工作流与依赖

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_typeengine_definition_idengine_revisiondeployment_metadata,回填现有记录为 n8n
  • 迁移期保留 n8n_workflow_idn8n_workflow_name,禁止破坏性删除。
  • 建立 workflow_schedulesworkflow_runsworkflow_task_runsworkflow_engine_bindingsworkflow_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

工作项:

  • 定义 validatedeploy_disabledactivatedeactivateexecutepausekillreplaybackfillget_executionget_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_uidpurpose、参数和 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/<date>_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 顺序交付与增量验证计划