十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

AutoGen 分布式运行时完整指南:gRPC 智能体跨节点协作快速上手

AutoGen 分布式运行时完整指南:gRPC 智能体跨节点协作快速上手 AutoGen 分布式运行时完整指南gRPC 智能体跨节点协作快速上手【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogenAutoGen 的分布式运行时基于 gRPC 提供多智能体跨节点协作能力一个中心 Host 负责消息路由多个 Worker 进程各自承载智能体Python 与 .NET 还能接入同一张消息网络。当你已经跑通单智能体示例、准备把 Agent 拆到不同进程甚至不同机器上时这套 Host/Worker 架构就是最短路径。单机跑多智能体为什么会卡住三个绕不过去的痛点单进程 InProcessRuntime 写 Demo 很顺但一旦认真做生产三个问题会接踵而至进程隔离没了退路。所有智能体挤在同一个 Python 进程里共享 GIL。一个智能体在做耗时的模型调用或代码执行时其他智能体的消息处理直接被拖住你也无法把模型推理节点和编排节点放到两台机器上各自扩缩容。跨语言是硬墙。AutoGen 同时有 Python 和 .NET 两套实现但 InProcessRuntime 只在单语言单进程内有效。C# 团队写的 Agent 和 Python 团队写的 Agent 之间没有任何进程内机制能搭桥。水平扩展无从下手。想给某个重度智能体比如跑检索、跑代码执行加副本或者让它独立重启而不影响其他智能体进程内运行时给不了这个粒度。结论先行出路是把消息传递从进程内存里搬出去换成一张显式的网络消息面。AutoGen 分布式运行时就是这张面而 gRPC 是它的传输协议。Host 与 Worker 是怎么把 gRPC 消息串起来的先给分工白话版GrpcWorkerAgentRuntimeHost一个 gRPC 服务端进程你可以把它理解为邮局路由表。它自己不跑业务逻辑只维护一张订阅表谁订了哪个主题、哪些智能体类型注册过然后把收到的消息投给订阅了对应 Topic 的 Worker。GrpcWorkerAgentRuntime跑在每个 Worker 进程里的客户端运行时。它和 Host 之间是一条OpenChannel双向流协议定义见 agent_worker.proto智能体的注册、订阅、发布都走这条流。TopicTopicId消息的寻址单元由type和source两段组成比如TopicId(data_pipeline, ingest)。发布publish是广播语义所有订阅了该主题的实例都会收到点对点请求则用send_messageRPC带响应。一条消息从发布到被消费的完整旅程如下时序两个容易混淆的点发布出去的消息在网络上被包成CloudEvent带 id、type、source 等属性而send_message点对点请求包成RpcRequest/RpcResponse——两者都在 gRPC 扩展实现 里有对应处理逻辑。最小可运行示例一条采集→分析→播报数据管道换一套新角色把闭环跑通CollectorAgent采集把读数发到ingest主题AnalystAgent分析消费后产出Insight发到digest主题BroadcasterAgent播报消费digest主题并打印。按启动顺序分三步。第一步先起 Host验证端口通了这段在验证 gRPC 服务能独立存活且被信号优雅关闭import asyncio from autogen_ext.runtimes.grpc import GrpcWorkerAgentRuntimeHost async def main() - None: host GrpcWorkerAgentRuntimeHost(address0.0.0.0:50052) host.start() # 在后台任务里启动 gRPC 服务 await host.stop_when_signal() # CtrlC / SIGTERM 时优雅停止 if __name__ __main__: asyncio.run(main())第二步接入第一个 Worker验证能注册、能发这段验证 Worker 能连上 Host、智能体类型能注册、消息能进入指定 Topic此时还没有订阅者消息会被 Host 路由但无人消费from dataclasses import dataclass from autogen_core import ( DefaultSubscription, RoutedAgent, TopicId, try_get_known_serializers_for_type, ) from autogen_ext.runtimes.grpc import GrpcWorkerAgentRuntime dataclass class SensorReading: value: float class CollectorAgent(RoutedAgent): def __init__(self) - None: super().__init__(Collector) async def main() - None: runtime GrpcWorkerAgentRuntime(host_addresslocalhost:50052) await runtime.start() runtime.add_message_serializer(try_get_known_serializers_for_type(SensorReading)) await CollectorAgent.register(runtime, collector, CollectorAgent) await runtime.add_subscription(DefaultSubscription(agent_typecollector)) await runtime.publish_message( SensorReading(42.5), topic_idTopicId(data_pipeline, ingest) ) await runtime.stop_when_signal()注意register这一步会向 Host 上报智能体类型add_subscription则把collector 订阅默认主题写进 Host 的路由表。第三步分析智能体跨进程接手验证 Topic 路由这段验证消息如何从进程 A 的ingest主题被进程 B 的订阅者消费并沿管道往下游走dataclass class Insight: summary: str class AnalystAgent(RoutedAgent): def __init__(self) - None: super().__init__(Analyst) message_handler async def on_reading(self, message: SensorReading, ctx: MessageContext) - None: insight Insight(summaryf读数 {message.value} 已校验) await self.publish_message( insight, topic_idTopicId(data_pipeline, digest) ) async def main() - None: runtime GrpcWorkerAgentRuntime(host_addresslocalhost:50052) await runtime.start() for msg_type in (SensorReading, Insight): runtime.add_message_serializer(try_get_known_serializers_for_type(msg_type)) await AnalystAgent.register(runtime, analyst, AnalystAgent) await runtime.add_subscription( DefaultSubscription( topic_idTopicId(data_pipeline, ingest), agent_typeanalyst ) ) await runtime.stop_when_signal()这里的关键是DefaultSubscription(topic_id..., agent_type...)不是订默认主题而是精确订到data_pipeline/ingest这才是跨进程消息流转的钩子。第四步播报收尾闭环验证这段验证最后一段订阅链看到控制台打印出播报内容即说明三进程协作闭环成立class BroadcasterAgent(RoutedAgent): def __init__(self) - None: super().__init__(Broadcaster) message_handler async def on_insight(self, message: Insight, ctx: MessageContext) - None: print(f[播报] {message.summary}) async def main() - None: runtime GrpcWorkerAgentRuntime(host_addresslocalhost:50052) await runtime.start() runtime.add_message_serializer(try_get_known_serializers_for_type(Insight)) await BroadcasterAgent.register(runtime, broadcaster, BroadcasterAgent) await runtime.add_subscription( DefaultSubscription( topic_idTopicId(data_pipeline, digest), agent_typebroadcaster ) ) await runtime.stop_when_signal()完整的多 Worker 演示可参考仓库里的 gRPC Worker 示例里面有 publish/subscribe 链和 RPC 两种模式的对照写法。从单机到多机分布式 Agent 部署走一遍把上面四个进程从一台机器的四个终端搬到多机器只有地址变了代码几乎不动。按顺序走统一一个环境变量。Host 监听地址和 Worker 的host_address都改为读环境变量避免端口写死在代码里# Host 进程 host GrpcWorkerAgentRuntimeHost( addressos.environ.get(PIPELINE_HOST, 0.0.0.0:50052) ) # 各 Worker 进程 runtime GrpcWorkerAgentRuntime(host_addressos.environ[PIPELINE_HOST])在机器 A 起 Host。PIPELINE_HOST0.0.0.0:50052 python run_pipeline_host.py确认 A 机器的 50052 端口对 B、C、D 机器可达。在 B 机起采集 WorkerPIPELINE_HOSTA机器IP:50052 python run_pipeline_collector.py观察 Host 日志里出现collector类型注册。在 C 机起分析 Worker、D 机起播报 Worker方式相同只是PIPELINE_HOST指向 A。验证闭环B 上触发一次SensorReading发布D 上应打印播报。不想开四个终端的话docker compose 等价表达如下services: pipeline-host: build: . command: python run_pipeline_host.py ports: [50052:50052] worker-collector: build: . environment: [PIPELINE_HOSTpipeline-host:50052] depends_on: [pipeline-host] worker-analyst: build: . environment: [PIPELINE_HOSTpipeline-host:50052] depends_on: [pipeline-host]两个部署要点其一示例代码走的是 insecure channeladd_insecure_portinsecure_channel机器之间不加密公网暴露前必须套 TLS 网关或内网隔离其二Host 是无状态的路由层Worker 才是干活层——扩容时加 Worker 进程即可不需要动 Host。最常踩的坑与调优点Worker 反复报UNAVAILABLE然后进程退出。原因Host 重启或网络抖动Worker 通道的内置重试策略对UNAVAILABLE最多 3 次、指数退避只覆盖短暂抖动连接彻底断开后读取循环不会自动重连。处理给 Host 配容器 restart 策略或 systemd 守护Worker 侧挂了就直接重启进程这是当前最省心的恢复路径。消息发出去了对面就是收不到。原因通常是三类订阅没加或topic_id拼错type/source 两段都要一致两端注册的序列化器类型名data_type不一致导致反序列化失败发布与订阅跑在不同的 Worker 而 Host 没查到路由。处理逐段核对TopicId、确认add_message_serializer两端都调过、在 Host 侧观察该类型的注册与订阅记录。大载荷下延迟偏高。原因payload 默认按 JSON 内容类型序列化。处理构造运行时传payload_serialization_formatPROTOBUF_DATA_CONTENT_TYPE直接走 protobuf 编码跨语言 Agent 本来就必须共享 protobuf schema协议定义在 protos/agent_worker.proto这条顺理成章。同一主题被多个同类型 Worker 重复处理。原因publish 是广播语义每个订阅实例都会收到一份。处理需要恰好一人处理的场景改用send_message点对点 RPC有请求/响应、request_id关联确实要多副本时按主题或 agent key 分片别让所有实例订同一个主题。想观测但不知道从哪看。GrpcWorkerAgentRuntime构造函数接受tracer_provider挂上 OpenTelemetry 后每次发布、处理都有 span跨进程调用链可以在 tracing 后端里串起来。下一步可以做什么 把AnalystAgent从规则判断换成 LLM 驱动用OpenAIChatCompletionClient等模型客户端驱动智能体分析环节直接调模型分布式骨架不用改。 接入 .NET 端跑一次真正的跨语言协作.NET gRPC 运行时 与 Python 端共用同一份 proto仓库里的 跨语言 Hello 示例 是最小参照。 试一下 RPC 模式与主题前缀订阅对照 官方 gRPC 示例 中的 RPC 变体把采集→分析改成请求/响应式调用。 把 OpenTelemetry 的tracer_provider接进所有 Worker先拿到一条完整的跨进程消息链路再谈扩容与告警。【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogen创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表