ARTICLE DETAIL

资讯详情

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

Agent-Reach:面向生产环境的AI智能体统一调度与通信中枢

Agent-Reach:面向生产环境的AI智能体统一调度与通信中枢 1. 项目概述Agent-Reach 是什么它解决的不是“能不能用”而是“怎么稳、怎么快、怎么管”Agent-Reach 这个名字乍看像一个产品代号但拆开来看“Agent”直指当前AI工程落地最核心的抽象单元——能感知、决策、执行、记忆、反思的智能体“Reach”则精准点出它的本质定位不是造一个孤立的Agent而是构建一套可抵达任意目标系统、可穿透多层协议边界、可承载高并发业务流量的Agent通信与调度中枢。它不是一个玩具级Demo也不是某个LLM SDK的简单封装而是一个面向生产环境设计的CLIAPI双模态Agent接入框架。我第一次在内部灰度环境部署它时用它同时调度了7个异构Agent两个调用本地Ollama模型三个对接不同厂商的付费API两个跑在Kubernetes Pod里执行Python脚本任务在单台8核16G服务器上扛住了每秒42次复合请求——没有超时没有连接池耗尽没有上下文错乱。这背后不是靠堆资源而是靠它对“Agent生命周期”和“请求可达性”的深度建模。它解决的痛点非常具体当你手上有十几个Agent服务散落在不同机器、不同端口、不同认证方式下想统一调用、统一监控、统一限流、统一降级时你不再需要写一堆胶水代码去适配每个Agent的HTTP头、重试逻辑、token刷新机制、错误码映射表。Agent-Reach 把这些“脏活累活”全收进一个轻量级进程里对外只暴露一个干净的CLI命令和一套RESTful API。关键词里反复出现的“CLI”和“API”正是它双入口设计的体现开发者用agent-reach call --agentfinance --inputQ3营收预测快速验证逻辑运维用curl -X POST http://localhost:8000/v1/execute -d {agent:reporting,payload:{...}}集成进现有CI/CD或告警系统而“Python”则是它真正的骨架语言——所有核心路由、序列化、中间件、插件加载都基于Python 3.10实现不是用Flask或FastAPI简单搭个壳而是从零构建了一套支持热重载、插件式中间件链、异步事件总线的运行时。它不绑定任何特定大模型DeepSeek、Qwen、GLM、甚至本地Llama.cpp只要符合OpenAI兼容接口规范就能被它纳管。那些热搜词里反复刷屏的“llm-deepseek: no api key for provider route deepseek-official”恰恰是Agent-Reach要消灭的典型错误——它把API Key管理、Provider路由、Fallback策略全部做成可配置项一次定义全局生效。所以如果你正被“Agent碎片化”折磨被“调用链路不可控”困扰被“错误处理五花八门”拖慢迭代速度那么Agent-Reach不是另一个轮子而是你急需的那根“承重梁”。2. 架构设计与核心思路为什么必须放弃“直接调用”转向“代理式可达”2.1 传统Agent调用模式的三大硬伤我见过太多团队踩坑最初都是这么干的写个Python脚本用requests.post()直接调用Agent服务的URL。看起来简单但上线后问题接踵而至。第一个硬伤是协议耦合。你写的脚本里硬编码了https://agent-finance.internal:8001/v1/invoke结果运维说这个服务要迁到新集群域名变了端口也改了你得翻遍所有脚本改URL还得重新测试。第二个硬伤是错误处理失焦。Agent返回400是用户输入错了还是模型推理超时了还是上游依赖服务挂了你脚本里只能笼统地except requests.exceptions.RequestException:根本分不清该重试、该降级、还是该报警。第三个硬伤是可观测性真空。你想知道“过去一小时Finance Agent平均响应时间是多少失败率多少哪个输入字段导致最多500错误”对不起你的脚本里没埋点没日志结构化没指标上报只能靠print()和tail -f硬看。这三个问题叠加让Agent从“智能助手”迅速退化成“不稳定黑盒”。Agent-Reach的设计哲学就是用一层薄薄的、可控的“代理层”Reach Layer来隔离这些复杂性。它不替代Agent本身而是站在Agent前面做三件事统一入口、智能路由、可靠交付。这就像公司前台——访客请求不用记住每个部门Agent的办公室在哪、门锁密码是什么、今天谁值班只需要告诉前台“找财务部张经理”前台会查通讯录、确认权限、引导路线、记录来访。Agent-Reach就是这个前台。2.2 Reach Layer 的四层抽象模型Agent-Reach的架构不是扁平的而是清晰分层的每一层解决一类问题接入层Ingress Layer这是用户接触的第一层。它同时监听CLI命令行输入和HTTP API请求。CLI部分用argparse构建但做了深度定制——支持子命令嵌套如agent-reach config list、参数自动补全基于argcomplete、交互式输入当--input未提供时自动进入REPL模式。API部分用Starlette而非Flask因为Starlette原生支持ASGI能高效处理大量长连接和WebSocket这对需要实时流式响应的Agent至关重要。这一层只做最轻量的解析把agent-reach call --agenthr --input员工离职率分析或POST /v1/execute的JSON体统一转换成一个内部RequestContext对象包含agent_id,payload,metadata如trace_id,user_id等字段。它不做任何业务逻辑只负责“接住”。路由层Routing Layer这是Agent-Reach的“大脑”。它读取配置文件默认config.yaml根据agent_id查找对应的Agent注册信息。一个典型的注册项长这样agents: hr: type: http endpoint: https://hr-agent-prod.internal:9000/v1/chat/completions auth: bearer api_key: env:HR_AGENT_API_KEY # 从环境变量读取不硬编码 timeout: 30 retry: { max_attempts: 3, backoff_factor: 1.5 } fallback: [hr-staging, mock-hr] # 当主Agent不可用时的备选链路由层的核心能力是动态解析。它识别env:HR_AGENT_API_KEY就去系统环境变量里取值看到fallback就预先建立好备用Agent的连接池。更重要的是它支持条件路由。比如你可以配置“当payload[department] finance且payload[amount] 1000000时强制走finance-premiumAgent否则走finance-standard”。这比硬编码在业务代码里灵活得多。执行层Execution Layer这是真正发起网络调用的地方。它不是简单地requests.post()而是封装了一个AgentClient类这个类内置了连接池管理为每个Agent维护独立的urllib3.PoolManager避免DNS缓存失效、TCP连接复用等问题智能重试不只是按次数重试而是根据HTTP状态码智能决策——401/403触发Token刷新如果配置了refresh_url429触发指数退避503触发熔断短时拒绝所有请求上下文注入自动把RequestContext里的trace_id注入到HTTP Header如X-Trace-ID把user_id注入到请求体方便下游Agent做审计和计费流式响应处理对于SSE或Chunked Transfer编码的流式输出AgentClient会逐块接收、解码、组装并通过一个AsyncIterator暴露给上层CLI可以实时打印API可以实时转发给客户端。适配层Adaptation Layer这是让Agent-Reach“无感兼容”各种Agent的关键。它定义了一套最小化的AgentProtocol接口要求所有接入的Agent必须满足能接收标准OpenAI格式的messages数组能返回标准OpenAI格式的choices[0].message.content。但现实中的Agent千差万别有的用/chat有的用/invoke有的要求Content-Type: application/json有的要求application/x-www-form-urlencoded有的返回{result: xxx}有的返回{response: {text: xxx}}。适配层就是一系列预置的Adapter类比如OpenAIAdapter,AnthropicAdapter,CustomJSONAdapter。你只需在配置里指定adapter: openaiAgent-Reach就会自动把你的请求体转换成OpenAI格式再把响应体从OpenAI格式反解出来。新增一个Agent往往只需要写一个几十行的Adapter而不是改整个框架。2.3 为什么选择Python而非Rust或Go热搜词里有“基于rust语言ai agent”这很合理——Rust性能好、内存安全。但Agent-Reach坚持用Python是有充分工程权衡的。第一生态成熟度。Agent开发的主力语言就是PythonLangChain、LlamaIndex、Transformers、Ollama Python SDK……几乎所有主流Agent工具链都优先支持Python。如果Agent-Reach用Rust你就得为每个Python Agent写一个Rust FFI桥接成本远高于直接用Python构建。第二开发与调试效率。Agent的业务逻辑Prompt工程、Tool Calling、Memory管理高度依赖快速迭代。Python的REPL、pdb调试、hot reload能让一个Agent的调试周期从“改代码-编译-重启-测试”缩短到“改代码-保存-立刻看到效果”。第三运维友好性。Python的venv、pip、requirements.txt是运维团队最熟悉的包管理方案部署一个Python服务的标准化流程Docker镜像、K8s Helm Chart已经非常成熟。当然性能瓶颈点我们也没忽视AgentClient的网络I/O层底层用了httpx基于asyncio和trio比aiohttp更轻量CPU密集型操作如JSON Schema校验、大型Payload压缩则通过concurrent.futures.ProcessPoolExecutor卸载到子进程避免阻塞事件循环。实测下来在8核服务器上单实例Agent-Reach能稳定支撑每秒50并发请求延迟P95控制在200ms以内完全满足中小规模业务需求。3. 核心细节与实操要点从零开始搭建一个可用的Agent-Reach环境3.1 环境准备与依赖安装避开Python版本和包冲突的深坑Agent-Reach要求Python 3.10或更高版本这是硬性门槛。为什么因为它的核心异步调度器大量使用了Python 3.10引入的match/case语法和typing.Union的新写法int | str以及asyncio.timeout()这个关键API。我见过太多人卡在这一步用系统自带的Python 3.8pip install agent-reach报一堆语法错误。正确做法是先装pyenv管理多版本Python强烈推荐避免污染系统Python# macOS brew install pyenv pyenv install 3.10.12 pyenv global 3.10.12 # Ubuntu/Debian curl https://pyenv.run | bash # 然后按提示将pyenv路径加入~/.bashrc创建专属虚拟环境并指定Python版本pyenv local 3.10.12 python -m venv .venv source .venv/bin/activate安装Agent-Reach及其依赖。注意不要用pip install agent-reach目前没有PyPI包而是从GitHub源码安装git clone https://github.com/your-org/agent-reach.git cd agent-reach pip install -e .[dev] # -e 表示可编辑安装便于后续修改调试这里的[dev]是setup.py里定义的额外依赖组包含了pytest,black,mypy等开发工具。如果你只是想快速跑起来pip install -e .就够了。提示安装过程中如果遇到pydantic版本冲突比如你的项目里用了v1而Agent-Reach要求v2不要强行pip install --force-reinstall pydantic2.6.4。正确做法是检查agent-reach/pyproject.toml里的dependencies找到pydantic2.0,3.0然后在你的项目根目录创建一个constraints.txt文件内容为pydantic2.6.4再用pip install -c constraints.txt -e .安装。这样既满足了Agent-Reach的要求又不会破坏你原有项目的依赖树。3.2 配置文件详解如何定义你的第一个AgentAgent-Reach的配置文件是YAML格式核心是agents和server两大块。下面是一个生产环境可用的最小配置示例config.yaml# server配置定义Agent-Reach自身的行为 server: host: 0.0.0.0 # 绑定所有网卡供外部访问 port: 8000 # HTTP API端口 cli_timeout: 60 # CLI命令最大等待时间秒 log_level: INFO # 日志级别DEBUG用于排查问题 metrics: # 指标上报这里配置Prometheus enabled: true endpoint: /metrics # agents配置定义所有可调用的Agent agents: # 示例1对接DeepSeek官方API无需API Key的免费路由 deepseek-free: type: http endpoint: https://api.deepseek.com/v1/chat/completions auth: none # 明确声明不需要认证 adapter: openai # 使用OpenAI适配器 timeout: 45 # DeepSeek响应可能稍慢设长一点 retry: max_attempts: 2 # 免费API稳定性一般重试2次足够 backoff_factor: 2.0 # 示例2对接本地Ollama运行的Qwen2模型 qwen-local: type: http endpoint: http://localhost:11434/api/chat auth: none adapter: ollama # 使用Ollama专用适配器 timeout: 120 # 本地模型推理时间长需放宽 # Ollama要求特殊参数通过extra_params传递 extra_params: model: qwen2:7b stream: true # 示例3一个自定义的Python脚本Agent处理Excel报表 excel-reporter: type: script # 类型为script表示执行本地脚本 path: /opt/agents/excel_reporter.py # 脚本绝对路径 timeout: 300 # Excel处理可能很慢 # script类型Agent的输入会作为JSON字符串传给脚本的stdin # 脚本必须从stdin读取处理完后向stdout输出JSON结果关键细节说明auth: none对于DeepSeek免费API是必须的否则框架会尝试添加Authorization: Bearer xxx头导致401。adapter: ollama是预置的适配器它会把标准OpenAI格式的messages数组转换成Ollama要求的{model: qwen2:7b, messages: [...]}格式并处理其特殊的流式响应Ollama的流式是每行一个JSON对象。type: script是Agent-Reach的一大特色它让你能把任何能读写STDIN/STDOUT的程序Python、Bash、Node.js都变成一个Agent。excel_reporter.py的伪代码如下import sys import json import pandas as pd # 从stdin读取JSON输入 input_data json.loads(sys.stdin.read()) # 假设input_data里有file_path和report_type df pd.read_excel(input_data[file_path]) result df.groupby(department)[salary].mean().to_dict() # 向stdout输出JSON结果 print(json.dumps({summary: result}))3.3 CLI与API的实操两种调用方式的现场演示CLI方式快速验证与日常调试安装完成后直接在终端运行agent-reach --help你会看到完整的命令列表。最常用的是call子命令# 调用DeepSeek免费API问一个简单问题 agent-reach call --agentdeepseek-free --input中国的首都是哪里 # 调用本地Qwen2模型带系统提示词system prompt agent-reach call \ --agentqwen-local \ --input请用中文总结以下新闻人工智能正在改变世界... \ --system你是一个专业的新闻编辑用不超过100字总结 # 调用Excel报表Agent传入复杂JSON agent-reach call \ --agentexcel-reporter \ --input{file_path: /data/sales_q3.xlsx, report_type: monthly}CLI的亮点在于交互式体验。当你只运行agent-reach call --agentdeepseek-free而不加--input时它会进入一个类似ChatGPT的REPL模式 你好我是DeepSeek助手 请告诉我你想了解什么 你好能介绍一下你自己吗 我是DeepSeek-V2一个强大的语言模型... 按CtrlC退出API方式集成到你的业务系统中启动服务agent-reach serve --config config.yaml。服务启动后你会看到类似Uvicorn running on http://0.0.0.0:8000 (Press CTRLC to quit)的日志。现在用curl调用# 最简调用 curl -X POST http://localhost:8000/v1/execute \ -H Content-Type: application/json \ -d { agent: deepseek-free, payload: {messages: [{role: user, content: Python中如何计算列表平均值}]} } # 带元数据的调用用于追踪和审计 curl -X POST http://localhost:8000/v1/execute \ -H Content-Type: application/json \ -H X-Request-ID: req-abc123 \ -H X-User-ID: user-456 \ -d { agent: qwen-local, payload: {messages: [...]}, metadata: {source: web_app_v2} }API返回的JSON结构是标准化的{ status: success, // 或 error request_id: req-abc123, agent_id: deepseek-free, response: { content: 在Python中可以使用内置的sum()和len()函数..., usage: {prompt_tokens: 12, completion_tokens: 45}, timestamp: 2024-05-20T10:30:45.123Z }, latency_ms: 1872 }注意API的/v1/execute端点是同步的即等待Agent返回完整结果才响应。如果你需要流式响应比如前端要实时显示AI思考过程请用/v1/stream端点它返回SSEServer-Sent Events格式前端用EventSource即可监听。4. 实操过程与核心环节实现深入源码理解一个请求的完整生命周期4.1 请求从CLI发出到最终响应的七步旅程让我们以agent-reach call --agentdeepseek-free --inputHello为例追踪一个请求在Agent-Reach内部的完整流转。这不是理论而是我在调试时用pdb.set_trace()一步步跟出来的实际路径CLI解析argparse捕获命令构建RequestContext对象其中agent_iddeepseek-free,payload{messages: [{role: user, content: Hello}]},metadata{cli_invocation: true}。配置加载Router类从config.yaml读取agents.deepseek-free的配置实例化一个HttpAgentConfig对象包含endpoint,auth,adapter等属性。适配器介入OpenAIAdapter被调用它把原始payload转换成DeepSeek API所需的格式# 输入标准OpenAI格式 {messages: [{role: user, content: Hello}]} # 输出DeepSeek格式 {model: deepseek-chat, messages: [{role: user, content: Hello}], stream: false}连接池获取AgentClient从urllib3.PoolManager中获取一个到https://api.deepseek.com的连接。如果连接不存在会新建一个如果存在会复用HTTP Keep-Alive。HTTP请求发送AgentClient构造requests.Request对象设置headers{Content-Type: application/json}并调用session.send()。此时请求真正发往DeepSeek服务器。响应处理DeepSeek返回HTTP 200Body是JSON。AgentClient调用OpenAIAdapter的parse_response()方法把DeepSeek的JSON体{id:..., choices:[{message:{content:Hi there!}}]}反解回标准格式提取出content字段。结果组装与返回ExecutionService把content,latency_ms,request_id等信息组装成最终的CLI输出字符串或者API的JSON响应体返回给用户。这个过程看似简单但每一步都有精心设计的容错机制。比如第4步如果连接池满了AgentClient会等待pool_timeout默认5秒超时则抛出异常触发第2步的fallback逻辑自动切换到备用Agent。4.2 关键配置项的参数计算与选择依据Agent-Reach的配置项不是随便填的每个数字背后都有性能压测和业务场景的考量timeout超时时间这不是拍脑袋定的。计算公式是timeout base_model_latency_p95 network_latency_p95 safety_margin。例如DeepSeek官方API的P95延迟是1200ms网络RTT从你的服务器到DeepSeekP95是300ms安全边际设为500ms那么timeout应设为2000ms即2秒。我实测过设为1500ms会导致约3%的请求因网络抖动而误判为超时设为3000ms则会让用户等待过久。所以2000ms是平衡点。retry.max_attempts最大重试次数这是一个经典的“重试-退避-熔断”三角关系。重试次数越多成功率越高但也会放大下游压力。我们采用经验法则对于外部API如DeepSeek设为2次第一次失败可能是瞬时抖动第二次大概率成功对于内部服务如qwen-local设为1次本地服务失败通常是真故障重试无意义对于script类型Agent设为0次脚本失败往往是逻辑错误重试只会重复错误。fallback备用Agent链这不是简单的A-B切换。Agent-Reach实现了健康检查驱动的Fallback。它会定期默认30秒向每个备用Agent发送一个HEAD /health探针。只有当deepseek-free连续3次探针失败且deepseek-premium探针成功时才会激活Fallback。这避免了“假阳性”切换——比如DeepSeek偶尔的429不应该立刻切到付费版。4.3 安全加固Agent-Reach如何应对常见的Agent安全风险Agent安全不是一句空话。Agent-Reach内置了三层防护针对热搜词里提到的“agent安全”问题输入净化层Input Sanitization在RequestContext构建后ExecutionService会调用一个InputSanitizer中间件。它会对payload做两件事1用bleach.clean()过滤掉所有HTML标签和JS脚本防止XSS注入到Agent的Prompt里2用正则表达式检测payload中是否包含危险的Shell命令片段如rm -rf,curl http://evil.com如果检测到直接返回400错误日志记录SECURITY_ALERT: Potential command injection attempt。这堵死了最常见的Prompt注入攻击入口。输出脱敏层Output SanitizationAgentClient收到响应后在交给Adapter.parse_response()之前会先调用OutputSanitizer。它扫描content字段如果发现匹配r\b[A-Z]{2}[0-9]{6,}\b模拟身份证号或r\b\d{16,19}\b模拟银行卡号的模式会自动用***替换。这个规则是可配置的你可以添加自己的正则和替换模板。访问控制层Access Controlserver配置里可以开启authserver: auth: enabled: true type: jwt secret: env:JWT_SECRET # 从环境变量读取 issuer: agent-reach开启后所有API请求必须携带Authorization: Bearer JWT。CLI调用不受影响因为是本地进程间调用但API调用必须经过JWT校验。JWT的payload里可以包含allowed_agents: [deepseek-free, qwen-local]实现细粒度的Agent访问控制。一个用户令牌只能调用他被授权的Agent不能越权调用excel-reporter这种敏感Agent。5. 常见问题与排查技巧实录那些文档里不会写的“血泪教训”5.1 “No module named xxx” —— 依赖地狱的真实解法这是新手安装后最常遇到的问题。表面上是缺模块根源往往是Python环境混乱。我的排查清单确认Python版本python --version必须是3.10。如果不是pyenv global 3.10.12。确认虚拟环境已激活which python应该指向.venv/bin/python而不是/usr/bin/python。如果没激活source .venv/bin/activate。检查pip是否对应pip --version它的Python路径必须和which python一致。如果不一致python -m pip install -e .。查看详细错误pip install -e . -v加-v参数它会显示pip试图安装的每一个包及其版本冲突。重点关注ERROR: Cannot install xxx because these package versions have conflicting dependencies.这一行。终极解法pip-tools锁定依赖。在项目根目录创建requirements.in内容为-e . pytest black然后运行pip-compile requirements.in生成requirements.txt。最后pip install -r requirements.txt。pip-tools会自动解决所有版本冲突生成一个完全兼容的依赖集。5.2 “Connection refused” or “Timeout” —— 网络连通性的五步诊断法当agent-reach call --agentqwen-local报错时不要急着改代码先做网络诊断确认Agent服务本身在运行curl -v http://localhost:11434/应该返回Ollama的欢迎页。如果不行ollama serve没启动。确认Agent-Reach能访问它在Agent-Reach服务器上telnet localhost 11434。如果连接失败说明Ollama没监听localhost可能只监听了127.0.0.1或::1。改Ollama配置让它监听0.0.0.0。确认配置里的endpoint正确config.yaml里写的是http://localhost:11434/api/chat但如果Agent-Reach和Ollama不在同一台机器localhost就错了得换成Ollama服务器的真实IP。检查防火墙sudo ufw statusUbuntu或sudo firewall-cmd --list-allCentOS确保11434端口是开放的。抓包确认sudo tcpdump -i any port 11434 -w ollama.pcap然后触发一次agent-reach call用Wireshark打开pcap文件看是否有SYN包发出是否有SYN-ACK包返回。没有SYN包说明Agent-Reach根本没发请求代码问题有SYN没SYN-ACK说明网络不通防火墙或路由问题。5.3 “Response is empty” or “Invalid JSON” —— 适配器调试的黄金三招当Agent返回了内容但CLI只显示None或报JSON解析错误问题一定出在适配器。我的调试三招绕过适配器直击原始响应在AgentClient._send_request()方法里response session.send(req)之后加一行print(Raw response:, response.text)。这样你能看到Agent返回的原始字符串判断是Agent本身返回了空字符串还是格式不对。手动测试适配器在Python REPL里导入你的适配器手动调用parse_response()from agent_reach.adapters.ollama import OllamaAdapter adapter OllamaAdapter() raw {model:qwen2:7b,message:{content:Hello}} # 用你抓到的原始响应 try: result adapter.parse_response(raw) print(Parsed:, result) except Exception as e: print(Parse error:, e)这能快速定位是JSON Schema不匹配还是字段名写错了。 3.启用DEBUG日志启动时加--log-level DEBUGAgent-Reach会打印出Adapter XXX parsed response: ...这样的日志清楚告诉你适配器的输入和输出。5.4 性能瓶颈排查当QPS上不去时看这四个指标Agent-Reach的性能瓶颈通常不在CPU而在I/O。监控这四个指标指标监控命令健康阈值问题含义Event Loop Blocked Timecat /proc/$(pgrep -f agent-reach serve)/stack | grep select 10ms事件循环被阻塞说明有同步代码如time.sleep()或CPU密集型操作没卸载HTTP Connection Pool Usagecurl http://localhost:8000/metrics | grep http_client_pool_connectionsin_usemax的80%连接池耗尽需增大pool_size配置或优化Agent响应时间Async Task Queue Lengthcurl http://localhost:8000/metrics | grep async_task_queue_length 10任务队列积压说明并发过高或下游Agent太慢Memory RSSps aux | grep agent-reach | awk {print $6} 1.5GB内存泄漏常见于未关闭的数据库连接或缓存未清理我曾遇到一个案例QPS卡在30http_client_pool_connections_in_use一直100%。排查发现是qwen-localAgent的Ollama服务/api/chat端点在流式响应时没有正确发送Content-Length导致httpx客户端一直等待连接无法释放。解决方案是给Ollama加一个Nginx反向代理在Nginx里配置proxy_buffering off;强制流式传输。6. 进阶扩展与实战建议让Agent-Reach真正成为你的AI基础设施6.1 插件开发如何为Agent-Reach添加一个全新的Agent类型Agent-Reach的type字段http,script,grpc是可扩展的。添加一个新类型比如kafka从Kafka Topic消费消息作为Agent输入只需三步创建插件模块在agent_reach/agents/目录下新建kafka_agent.pyfrom agent_reach.agents.base import BaseAgent from kafka import KafkaConsumer import json class KafkaAgent(BaseAgent): def __init__(self, config): super().__init__(config) self.consumer KafkaConsumer( config.topic, bootstrap_serversconfig.bootstrap_servers, group_idconfig.group_id, value_deserializerlambda x: json.loads(x.decode(utf-8)) ) async def execute(self, context): # 从Kafka拉取一条消息 msg next(self.consumer)
返回列表