工作流编排
概述
工作流编排平台提供了一套完善的能力,用于管理复杂的多步骤流程——这些流程往往需要在多个系统、服务和智能体之间进行协调。此类平台着重关注可靠性、可扩展性,以及企业级的编排能力。
主流编排平台
Littlehorse
平台:Littlehorse.io
架构:微服务、工作流、集成与 AI 智能体
技术:基于开源 LittleHorse Kernel 构建
核心特性
微服务编排: - 原生支持微服务架构模式 - 服务网格集成能力 - 分布式事务管理 - 容错与弹性模式
工作流管理: - 可视化工作流设计器与管理界面 - 支持分支和并行执行的复杂工作流编排 - 跨工作流步骤的状态管理与持久化 - 实时工作流监控与调试
AI 智能体集成: - 原生支持 AI 智能体编排 - 智能体间通信协议 - 基于智能体能力的智能工作流路由 - 与主流 AI 框架和平台的集成
企业级集成: - 企业服务总线(ESB)能力 - 遗留系统集成支持 - API 网关与管理功能 - 事件驱动架构支持
架构组件
LittleHorse Kernel: - 工作流引擎:核心编排与执行引擎 - 状态管理:跨服务的分布式状态管理 - 事件处理:实时事件处理与路由 - 服务注册中心:动态服务发现与注册
集成层: - API 网关:集中式 API 管理与路由 - 消息代理:服务间的可靠消息传递 - 协议适配器:支持多种通信协议 - 数据转换:内置数据映射与转换
管理与监控: - 仪表板:基于 Web 的管理与监控界面 - 分析:工作流性能分析与洞察 - 告警:实时告警与通知系统 - 审计追踪:完整的审计日志与合规支持
应用场景
企业级集成: - 遗留系统现代化与集成 - 多系统数据同步 - 跨部门业务流程自动化 - API 编排与管理
微服务协调: - 服务间通信编排 - 分布式事务管理 - 熔断器与重试模式实现 - 服务网格集成与管理
AI 智能体工作流: - 多智能体系统协调 - 智能体任务分发与负载均衡 - 基于智能体能力的智能工作流路由 - 智能体性能监控与优化
DevOps 与自动化: - CI/CD 流水线编排 - 基础设施自动化与资源供给 - 部署工作流管理 - 监控与告警自动化
技术能力
可扩展性与性能: - 跨节点水平扩展 - 高吞吐量消息处理 - 低延迟工作流执行 - 基于工作负载的自动弹性伸缩
可靠性与容错: - 内置重试机制与熔断器 - 状态管理的分布式共识 - 自动故障转移与恢复 - 数据一致性与完整性保障
安全与合规: - 基于角色的访问控制与权限管理 - 传输中和静态数据加密 - 审计日志与合规报告 - 与企业身份系统集成
快速开始
安装与配置:
# Docker installation
docker run -d --name littlehorse \
-p 8080:8080 \
littlehorse/littlehorse:latest
# Kubernetes deployment
kubectl apply -f https://raw.githubusercontent.com/littlehorse-enterprises/littlehorse/main/k8s/deployment.yaml
基础配置:
# littlehorse.yaml
server:
port: 8080
host: 0.0.0.0
database:
type: postgresql
host: localhost
port: 5432
name: littlehorse
messaging:
broker: kafka
bootstrap_servers: localhost:9092
开发工作流: 1. 设计工作流:使用可视化设计器创建工作流定义 2. 定义服务:注册微服务及其能力 3. 配置集成:建立与外部系统的连接 4. 部署与测试:部署工作流并用示例数据进行测试 5. 监控与优化:使用监控工具优化性能
集成模式
事件驱动架构: - 事件溯源与 CQRS 模式支持 - 实时事件流处理 - 事件关联与聚合 - 分布式事务的 Saga 模式
API 编排: - RESTful API 组合与编排 - GraphQL 联邦与模式拼接 - 限流与节流 - API 版本控制与生命周期管理
数据流水线编排: - ETL/ELT 工作流编排 - 实时数据流流水线 - 数据质量与校验工作流 - 多数据源集成
最佳实践
工作流设计: 1. 模块化设计:创建可复用的工作流组件 2. 错误处理:实施全面的错误处理策略 3. 状态管理:设计具有适当持久化的有状态工作流 4. 性能优化:针对吞吐量和延迟优化工作流 5. 测试策略:对复杂工作流实施充分的测试
卓越运营: 1. 监控:实施全面的监控与告警 2. 日志:结构化日志,用于调试与审计追踪 3. 安全:在整个平台中应用安全最佳实践 4. 备份与恢复:实施稳健的备份与灾难恢复方案 5. 容量规划:为增长与扩展需求提前规划
Temporal
网站:temporal.io
许可证:MIT(开源服务端 + SDK)
架构:持久化执行平台——基于代码的工作流(Workflows-as-Code),配合事件历史回放机制
Temporal 是什么
Temporal 是一个持久化执行平台,能让长时间运行的容错工作流像普通应用程序代码一样简单易写。官方文档将其描述为:"持久化执行通过保证应用程序运行至完成,确保其在不利条件下也能正确运行。"
开发者无需在应用逻辑中管理重试、状态检查点和故障恢复,而是直接编写工作流代码,就好像它从不中断地连续运行。一旦发生崩溃或中断,Temporal 会通过回放已记录的事件历史来恢复工作流状态,从中断点精确恢复执行。这使 Temporal 尤其适合调用不稳定 LLM API、运行时间长达数分钟乃至数小时、且需要在多个步骤间保持连贯状态的 AI 智能体流水线。
一个重要的架构说明:Temporal Service 负责编排并持久化状态,但运行你代码的是 Worker。Worker 向 Temporal Service 轮询任务,在你的基础设施中执行对应的工作流或活动代码并返回结果——你的数据始终不会离开你的掌控。
核心原语
工作流(Workflow):定义整体智能体逻辑的持久化有状态函数,包括活动序列、分支、信号和子工作流。必须是确定性的(Temporal 通过回放事件历史来实现崩溃恢复)。
活动(Activity):独立的工作单元——一次 LLM API 调用、工具调用、网络搜索或数据库读取。活动是工作流编排的非确定性、具有副作用的操作。每个活动都可配置重试策略和超时时间;活动失败时,Temporal 会根据配置自动重试。
Worker:长期运行的进程,向 Temporal Service 轮询任务并执行工作流/活动代码。Worker 是无状态的,支持水平扩展——增加更多 Worker 即可提升吞吐量,无需额外的协调开销。
事件历史(Event History):完整持久化的日志,记录工作流执行生命周期中的每一个事件,持久化到 Temporal Service 数据库。这是持久化执行的基础——Worker 崩溃后,通过回放事件历史可精确重建内存状态,并从故障点恢复执行,如同故障从未发生。
信号/更新(Signal / Update):发送到运行中工作流的异步消息(信号)或带响应的同步信号(更新)。原生支持人在回路(HITL)关口——暂停执行,等待人工审批、新数据输入或取消操作。
查询(Query):只读的同步操作,返回当前工作流状态而不影响执行。适用于暴露智能体进度或中间结果。
子工作流(Child Workflow):由父工作流派生的工作流,支持层级式多智能体架构。主管工作流派生专属子智能体工作流,收集结果后继续编排。
智能体 AI 模式
| 模式 | Temporal 的实现方式 |
|---|---|
| 容错 LLM 链 | 每次 LLM 调用是一个活动;瞬态 API 错误(限流、超时、5xx)自动以指数退避方式重试 |
| 长时间运行的研究智能体 | 工作流可运行数小时乃至数天;状态在基础设施故障期间持久化,无需修改代码 |
| 人在回路审批 | 工作流无限期暂停等待信号;无需轮询循环或外部状态存储 |
| 多智能体扇出/扇入 | 父工作流派生 N 个并行子工作流(各为一个子智能体),等待全部完成后聚合结果 |
| 智能体版本化部署 | workflow.get_version() API 将现有长时间运行的执行路由到旧代码,新执行路由到新代码——无需大规模迁移 |
| Saga / 补偿事务 | 每个步骤都有在失败时运行的补偿活动,在无需分布式事务协议的情况下保证最终一致性 |
SDK 支持
七种语言 SDK(多语言团队可在工作流和活动间混用不同语言):Go、Python、TypeScript/JavaScript、Java、.NET (C#)、PHP、Ruby。
部署选项
| 模式 | 说明 |
|---|---|
| Temporal Cloud | 全托管 SaaS;按用量计费;零运维开销 |
| 自托管(Docker Compose) | 本地开发/测试的单节点部署(docker-compose up) |
| 自托管(Kubernetes) | 通过 Helm chart 实现生产级部署 |
| 内嵌(Embedded) | 进程内 Temporal Service,用于单元测试和集成测试 |
最佳实践
| 挑战 | 解决方案 |
|---|---|
| 确定性 | 工作流代码中禁止使用随机数、当前时间或 I/O——使用 workflow.now(),并将所有副作用移至活动中 |
| 活动粒度 | 将每次 LLM 调用或外部 API 调用单独封装为一个活动;仅在整批操作具有原子性时才进行批量处理 |
| 负载大小 | 活动间传递引用(S3 key、数据库 ID)而非完整 LLM 响应;使用 Data Converter 进行压缩 |
| 超时调优 | 根据 P99 延迟设置 start_to_close_timeout;对长时间运行的活动使用心跳机制 |
| Worker 池 | 为 CPU 密集型(嵌入)、I/O 密集型(LLM 调用)和 GPU(本地推理)活动分别配置独立的 Worker 池 |
Confluent 数据流平台(Apache Kafka & Apache Flink)
网站:confluent.io
核心技术:Apache Kafka®(事件流,云原生引擎品牌为 Kora)与 Apache Flink®(流处理)
架构:事件驱动数据流平台,作为智能体框架之下的通信与数据层,而非独立的智能体编排器
Confluent 是什么
Confluent 将其数据流平台定位为智能体系统的"中枢神经系统"——并非 LangGraph、AutoGen、CrewAI 等智能体框架的替代品,而是这些框架中智能体发布和消费事件的事件骨干网络。该平台围绕四大支柱构建:
| 支柱 | 作用 |
|---|---|
| 流(Stream) | 通过云原生 Kafka 引擎 Kora 持续捕获并共享实时事件。 |
| 连接(Connect) | 通过 120+ 个预置和自定义连接器集成异构数据源,消除智能体与外部系统之间的硬编码依赖。 |
| 处理(Process) | 使用 Flink 流处理(SQL/Table API、JOIN、过滤器、Flink AI 模型推理)以实时上下文丰富数据,支持智能体 RAG 和实时嵌入流水线。 |
| 治理(Govern) | 流治理(Stream Governance)强制执行数据血缘、质量控制和可追溯性,确保智能体消费的数据安全且可验证。 |
智能体 AI 模式
| 模式 | Confluent 的实现方式 |
|---|---|
| 编排者-工作者(Orchestrator-Worker) | 编排者向 Kafka 主题发布带 key 的命令消息;工作者智能体以消费者组形式从分配的分区拉取消息,通过消费者再均衡协议和偏移量回放实现扩展与恢复 |
| 层级多智能体(Hierarchical multi-agent) | 将编排者-工作者模式递归应用;每一层将目标作为事件发布给下一层 |
| 黑板(Blackboard) | 以 Kafka 主题作为共享知识库;智能体通过发布/订阅更新,而非查询共享数据库 |
| 基于市场的协调(Market-Based coordination) | 竞价智能体将报价/请求作为事件发布;做市服务进行匹配,用中央事件日志取代 O(n²) 的点对点协商 |
| 智能体 RAG(Agentic RAG) | Flink 实时摄取并嵌入非结构化源数据,通过 Sink 连接器将向量存储到数据库(如 MongoDB Atlas),用于低延迟检索 |
| 故障恢复(Fault recovery) | Kafka 的日志架构允许智能体回放事件并从故障中恢复;幂等处理避免重试时的重复操作;死信队列处理需要人工介入的故障 |
完整的模式详解(传统方式与事件驱动方式的对比处理)见多智能体系统的事件驱动设计模式(Confluent)。
部署选项
| 模式 | 说明 |
|---|---|
| Confluent Cloud | 全托管 Kafka(Kora 引擎)+ Flink;按用量计费 |
| 自管理 Kafka/Flink | 开源 Apache Kafka 和 Apache Flink,自行托管 |
最佳实践
| 挑战 | 解决方案 |
|---|---|
| 工作者/智能体故障恢复 | 使用 Kafka 消费者组,让消费者再均衡协议自动重新分配负载;从最后保存的偏移量回放,而非构建自定义故障处理逻辑 |
| 重试时的重复副作用 | 将事件触发的智能体操作设计为幂等;将反复失败的事件路由至死信队列供人工审查 |
| 智能体决策的上下文过时 | 使用 Flink SQL/Table API 和 Flink AI 模型推理,在查询执行时以实时上下文丰富数据,而非依赖批量更新 |
| 共享事件流中的敏感数据 | 对事件负载应用字段级加密和访问控制;通过流治理强制实施血缘追踪和质量保障;设置符合法规(如 GDPR)的细粒度数据保留策略 |
| 工具/智能体集成蔓延 | 使用 120+ 预置连接器,并结合现有智能体框架(LangGraph、AutoGen、CrewAI)处理工具调用,避免硬编码点对点集成 |
注意事项
Confluent 是数据流和事件骨干网,而非智能体框架——必须与编排/智能体框架(如 LangGraph)配合使用,以实现推理和工具调用逻辑。最适合已在运行 Kafka/Flink 基础设施的组织,或以事件量、扇出和跨系统集成(而非持久化长时间运行的工作流状态,这是 Temporal 的专注点)为主要扩展挑战的多智能体系统构建场景。
Kestra
官网:kestra.io | 仓库:kestra-io/kestra 许可证:Apache 2.0(完全开源;企业版增加治理特性) 架构:事件驱动的编排平台,采用声明式 YAML 工作流;控制平面、执行器、调度器与工作器均为模块化组件 技术栈:Java 后端(UI 用 Vue.js/TypeScript);元数据后端支持 PostgreSQL、MySQL 或 H2;基于队列的任务分发以实现横向扩展
Kestra 是什么
Kestra 是一个开源、事件驱动的编排平台,面向数据、AI 与基础设施工作流,在一套声明式、语言无关的接口背后统一了定时与事件驱动两类自动化。工作流("flow")以 YAML 定义——可通过内置代码编辑器、API、Git 或无代码编辑器编写——并可在原生插件任务旁内嵌 Python、Node.js、R、Go 或 Shell 脚本。随着 Kestra 1.0(其首个长期支持 LTS 版本)发布,项目将自身重新定位为"声明式智能体编排平台(Declarative Agentic Orchestration Platform)",在生产强化特性(单元测试、SLA、插件版本化、Helm chart)之外,新增了原生 AI Agent 任务、一个 Copilot,以及一个模型上下文协议(MCP)服务器。
智能体 AI 能力
AI Agent 任务:启动一个由 LLM、记忆和工具(网页搜索、任务执行、调用其他 flow)驱动的自主进程,它动态决定采取哪些动作、以何种顺序执行,而非遵循固定的任务序列。记忆让智能体在多次执行间保留上下文,以指导后续提示。
MCP 集成:Kestra MCP 服务器把 flow 与执行管理(列出 flow、触发运行、管理命名空间文件)暴露给兼容 MCP 的 AI 工具(如 Claude Code、Cursor);Kestra 也可充当 MCP 客户端,让 AI Agent 任务调用外部 MCP 工具服务器。
智能体工作流的护栏:对高风险动作设人在回路审批步骤、对工具/LLM 瞬时故障设重试逻辑、用超时限制失控执行,以及为长时运行的智能体循环提供有状态/持久化执行——全部以声明式定义,并在 Kestra UI 中完全可观测。
多智能体编排:智能体既可独立运行,也可组合成多智能体系统(例如一个编排者 flow 调用若干专用子智能体 flow),同时始终作为代码保持可检视、可治理,而非智能体框架内部的黑盒。
核心特性
| 领域 | 能力 |
|---|---|
| 工作流定义 | 声明式 YAML,配合基于 Git 的版本控制;无代码编辑器与 REST API 作为备选编写方式 |
| 触发器 | 定时(cron)与实时事件驱动触发器——文件到达、消息总线事件(Kafka、Redis、Pulsar、AMQP、MQTT、NATS、AWS SQS、Google Pub/Sub、Azure Event Hubs) |
| 插件生态 | 1,500+ 集成/任务,覆盖数据库、云存储、API 与脚本语言 |
| 可靠性 | flow 层面内建重试、条件逻辑、动态任务与错误处理 |
| 规模 | 基于队列的工作器架构,设计上可扩展至数百万次工作流执行 |
部署方式
| 模式 | 说明 |
|---|---|
| 开源(自托管) | Apache 2.0;Docker、Docker Compose,或用于 Kubernetes 的 Helm chart(自 1.0 起稳定) |
| Kestra Cloud | 全托管 SaaS 服务 |
| 企业版 | 在开源核心之上增加 RBAC、审计日志、SSO、自助式"Apps"表单与工作器隔离 |
注意事项
Kestra 与 Apache Airflow(数据流水线编排)以及 Temporal(持久化、代码优先的工作流执行)有所重叠,但它以 YAML 优先/语言无关的编写模型、更广的开箱即用触发器与插件面,以及——自 1.0 起——一等公民的 AI Agent 任务与 MCP 连通性使自己区别开来,这些特性正对准智能体 AI 用例,而无需在其上另外拼接一个独立的智能体框架。
与其他编排方案的比较
企业级编排平台
Apache Airflow: - 优势:成熟平台、社区活跃、基于 Python - 应用场景:数据流水线编排、批处理 - 注意事项:主要面向数据工作流,不适合实时编排
Kubernetes: - 优势:容器编排、云原生、生态系统完善 - 应用场景:容器部署与管理、微服务编排 - 注意事项:以基础设施为核心,业务工作流编排需要额外工具支撑
Temporal: - 优势:通过事件历史回放实现持久化执行、高容错性、支持 7 种语言的友好 SDK API、通过信号原生支持人在回路、通过子工作流实现层级多智能体组合 - 应用场景:长时间运行的 AI 智能体工作流、容错 LLM 链、多智能体扇出/扇入、人工审批关口、智能体版本化部署 - 注意事项:Workflows-as-Code 方式需要具备编程能力;非无代码工具;工作流代码必须满足确定性以支持回放
Confluent(Apache Kafka & Apache Flink): - 优势:高扇出多智能体协调的事件驱动骨干(编排者-工作者、层级、黑板、基于市场等模式)、120+ 连接器用于系统集成、用于智能体 RAG 的实时流处理和嵌入流水线、流治理满足合规需求 - 应用场景:高并发事件驱动智能体流水线、基于市场/交易风格的智能体协调、实时 RAG 摄取、需要多个独立生产者/消费者松耦合的智能体生态系统 - 注意事项:本身不是智能体框架——必须与 LangGraph、AutoGen、CrewAI 等配合使用以实现智能体推理逻辑;针对事件吞吐量和扇出进行了优化,而非长时间运行的持久化工作流状态(Temporal 的专注点)
Kestra: - 优势:Apache 2.0 完全开源;声明式 YAML 编写,配合 Git 版本控制;1,500+ 插件;广泛的原生触发器支持(cron 与多种消息总线);自 1.0 起原生 AI Agent 任务,带记忆/工具;内建 MCP 服务器与 MCP 客户端支持 - 应用场景:数据 + AI + 基础设施的统一编排、事件驱动自动化、希望以声明式(YAML)而非代码优先的智能体框架来定义智能体工作流的团队 - 注意事项:在持久化智能体编排上比 Temporal 更年轻;AI Agent 任务与 MCP 支持仅在 1.0 LTS 版本才引入,因此相较于 Kestra 已经站稳的通用数据流水线编排,其智能体用例的生态成熟度仍在发展中
选型标准
技术要求: - 可扩展性:满足当前及未来工作负载需求的能力 - 可靠性:容错与恢复能力 - 性能:延迟与吞吐量要求 - 集成性:与现有系统和技术的兼容性
运营要求: - 管理:部署、配置和管理的便捷性 - 监控:可观测性与调试能力 - 安全:安全特性与合规支持 - 支持:厂商支持与社区资源
业务要求: - 成本:包括许可证和运营在内的总拥有成本 - 上市速度:实施与部署的速度 - 厂商锁定:可移植性与厂商独立性 - 未来路线图:平台演进与功能开发
参见
- Temporal — 面向智能体 AI 的持久化工作流编排
- Kestra — 声明式智能体编排平台
- 开源工作流引擎
- 标准 — 模型上下文协议(MCP)
- 多智能体系统
- 多智能体系统的事件驱动设计模式(Confluent)
- 生产环境部署
- 可观测性
参考资料
- Falconer, S. (2025). A Guide to Event-Driven Design for Agents and Multi-Agent Systems. Confluent, Inc. — Confluent 数据流平台小节的来源。
- Kestra GitHub Repository — 许可证、架构与特性概览
- Introducing Kestra 1.0: The Declarative Agentic Orchestration Platform — LTS 版本、AI Agent 任务、Copilot、MCP 服务器
- AI Tools in Kestra: Copilot, Agents, MCP Server & More — AI Agent 任务、记忆、工具与 MCP 集成细节
- AI Agents in Kestra – Autonomous Orchestration — AI Agent 任务机制、护栏、多智能体组合