DataOps Platform 当前通过 n8n 定义、调度和执行数据加工流程。n8n 使用带有
使用限制的 Sustainable Use License,且当前集成把平台业务模型绑定到了
n8n_workflow_id、n8n Workflow JSON 和 SSH 脚本执行方式。
平台已经具备稳定 DataFlow UID、环境级版本激活、RBAC、Outbox、外部数据源 加密凭据和 Worker 本地连接池。后续目标不是继续由人员在第三方调度 UI 中维护 流程,而是由调度智能体依据血缘、SLA、运行历史和资源状态进行规划,通过 MCP 完成受控发布、执行、恢复和优化。
Scheduling MCP Gateway,只暴露经过权限和策略约束的复合操作。Context MCP,向智能体提供 DataFlow、血缘、SLA、数据源能力、
连接池健康和历史运行指标,但不返回用户名、密码、密文或连接串。WorkflowSpec 和 SchedulePlan;确定性编译器将通过
校验的版本转换为 Kestra Flow。不得把未经校验的模型输出直接发布到生产。DataOps Runner 执行。Kestra 只传递 DataFlow/任务版本、
data_source_uid、用途和参数;Runner 复用现有
DataSourceConnectionManager,Kestra 不保存业务数据源凭据。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 |
智能体可以:
智能体不可以:
本方案只依赖 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 同时存在。详细实施工作包和退出门槛见 Kestra AI-first Data Factory 双轨改造实施计划。