
1. 先交代背景批量任务卡在HTTP上以后我怎么做内部服务通信我前阵子接手一个内部数据处理平台里面最核心的场景是批量文件分析。一批文件进来调度器要按顺序丢给不同分析节点处理每份文件本身不大但量很大最后要求整批跑完的时间越短越好。刚开始是用HTTP接口做的调度器拿到一批文件循环去请求各个分析服务的接口等返回结果再汇总。看起来没什么问题但是一跑起来就发现两个致命麻烦。第一是同步等待。一份文件的分析耗时要三到五秒调度器挨个请求一个节点慢一点后面的任务全部排队。加asyncio并行请求能救一部分但HTTP框架在处理这类内部服务通信时请求头、序列化、连接管理、重试逻辑都得自己反复造轮子代码越长维护成本越高。第二是协议疲劳。内部服务之间真正需要的是“调一个函数传参数拿返回值”。HTTP能实现但中间隔了一层URL路由、状态码、Content-Type、JSON序列化细节每加一个接口服务端要多写一段路由挂接客户端要多写一个封装函数。一两百个函数全铺开以后双份维护变成常态改一个字段名要牵动两端。我后来在项目里逐步把内部节点间的通信切到了RPC方案上选用的库是a2rpc名字里的“a2”我理解下来是Async To Async的意思也就是面向Python asyncio生态的远程过程调用组件。它解决的核心痛点是我不用再用HTTP那套繁琐的协议表达“调用”而是像本地函数一样调用远程方法连参数带返回值一起走异步端到端都照顾到了。想快速体验到这种便利可以参考a2rpc包提供的全部语法与参数再结合本文后面两次实际迁移记录来验证它到底适合什么场景。如果你现在也在搭内部微服务、爬虫集群分片调度、或异步任务编排又在犹豫要不要上gRPC或消息队列这篇文章值得你看完。我尽量少讲空话直接讲a2rpc怎么装、怎么用以及两个生产案例里它怎么落地的。不追求万金油方案只给真实场景里验证过的东西。2. 快速上手三分钟跑通一个最小的远程调用2.1 安装与版本说明安装很简单一条命令pip install a2rpc我这里实际用的是0.4.x版本自带的依赖很少核心就依赖pydantic和msgpack。注意一下早期0.1.x版本API差异比较大网上搜教程时最好直接看官方文档和Changelog以最新API为准。装完后可以确认一下版本python -c import a2rpc; print(a2rpc.__version__)2.2 服务端写一个可以被远程调用的方法a2rpc的入口概念和多数Python RPC库差不多定义一个普通的类里面的方法通过装饰器标记为可被远程调用然后把类的实例交给RPC服务器管理。# service.py import asyncio from a2rpc import RPCServer, rpc_method class CalcService: rpc_method async def add(self, a: int, b: int) - int: return a b rpc_method def get_name(self) - str: return calc-node async def main(): server RPCServer(CalcService(), bind0.0.0.0, port9100) await server.serve() if __name__ __main__: asyncio.run(main())跑起来python service.py像add这样标记了rpc_method的方法客户端可以直接当成本地函数调用。2.3 客户端远程调用长什么样# client.py import asyncio from a2rpc import RPCClient async def main(): async with RPCClient(127.0.0.1, 9100) as client: res await client.add(a1, b2) print(res) # 3 asyncio.run(main())第一次跑通这个例子时我最大的感受是写调用方代码的时候几乎不用关心网络层。add(a1,b2)这行代码和本地调用几乎没有区别。内部的连接建连、请求序列化、响应反序列化都被a2rpc封装掉了。有一点值得注意如果你的方法参数使用了纯关键字传参a2rpc底层的参数序列化是按关键字匹配的。我第一次没注意用位置参数传参远程方法签名比较长时容易对不上后来一律改用关键字形式清晰也更安全。3. 语法地图装饰器、注册方式与调用链拆解用a2rpc之前有必要先画出它的语法主干这样后面改代码时不会迷路。整个库的核心语法其实就三大块rpc_method装饰器、RPCServer的注册逻辑、RPCClient的调用链。3.1 rpc_method装饰器rpc_method可以装饰async函数也可以装饰普通函数。装饰async函数时a2rpc会把它并入事件循环调度装饰普通函数时框架会自动用asyncio.to_thread之类的机制转成异步执行避免阻塞事件循环。rpc_method(version2, timeout5) async def heavy_task(self, payload: dict) - dict: ...timeout参数单位是秒。这个参数非常关键它决定服务端在方法执行超时后直接返回错误还是继续等下去。默认情况下不同版本有差异我用的0.4.x版本默认是30秒生产环境建议显式配置。装饰器还可以传version用来做方法级版本管理。加了这个参数之后客户端可以指定调用某个具体版本的方法。内部服务升级过程中旧版本不一定全部立刻下掉这个参数就派上了用场。3.2 RPCServer注册流程一个RPCServer实例可以注册多个服务类server RPCServer(bind0.0.0.0, port9100) server.register(CalcService()) server.register(FileService()) server.register(DeviceService(), prefixdev)注意prefix参数如果设置了前缀客户端调用时方法名会变成dev_xxx_method这种形式。这个设计在多业务模块合用一个RPC服务端口时非常实用能在命名空间上隔离不同模块的同名方法。注册逻辑内部大致分两步遍历类的所有公共方法筛选带rpc_method标记的方法然后把这些方法名和函数对象映射到一张方法路由表。所以一个类里没被装饰的方法不会暴露给客户端。这里有个细节容易被忽略实例属性携带状态。服务端每次调用同一个远程方法时实例的self状态是持续存在的。也就是说你可以把一个计数器、连接池、缓存挂在实例上。比如class CounterService: def __init__(self): self._count 0 rpc_method def incr(self): self._count 1 return self._count这个特性用好了可以减少很多外部存储依赖但也要求你心里有数这个服务的生命周期是整个RPC服务进程的生命周期不是每次调用都新建实例。3.3 RPCClient调用链客户端的调用链相对简洁连接管理RPCClient支持作为异步上下文管理器使用也支持手动connect()和close()。请求编码调用一个方法时客户端会把方法名、参数、调用元信息打包成一条消息。序列化默认使用msgpack比JSON体积更小编解码速度也更快。响应解码服务端返回结果后再解码回来。client RPCClient(127.0.0.1, 9100) await client.connect() try: result await client.get_name() finally: await client.close()需要注意的是同一个连接上的多个请求是异步并发处理的不是串行阻塞。这意味着你可以同时发出多个请求而不必排队等待。这在批量调用场景中价值极大。后面的压测案例里我正是因为这一点把整批处理时间压缩到了原来的三分之一。4. 参数清单从启动到调用真正需要调出心得的参数就这几个a2rpc的参数不算多但每个都有实际意义。我按“启动参数、方法参数、客户端参数”三类整理了一张常用参数清单。4.1 服务端启动参数服务端启动时RPCServer支持这些参数参数含义默认值我的建议bind监听地址127.0.0.1生产环境按需设为0.0.0.0port监听端口无默认必须给选不常用的高位端口backlog底层监听队列长度100并发量高时适当调大max_workers同步方法转异步时的线程池大小默认CPU核心数同步方法多时调大serializer序列化方案msgpack二进制服务场景保持默认health_check_port健康检查HTTP端口无K8s探活建议开启auth_token鉴权令牌无跨网段时强烈建议配置health_check_port是独立于RPC端口的一个小HTTP服务用来返回OK字符串。K8s里做存活探针很顺手不会干扰到RPC繁忙时的健康检查。4.2 装饰器与方法参数rpc_method装饰器本身的参数我已经提到了version和timeout另外还有rate_limit和audit。rate_limit参数用来做方法级限流设置后同一方法在一秒内的最大调用次数。曾经我在做对外数据服务时就有第三方调用方会突然狂拉数据。给每个昂贵查询方法设置rate_limit后服务端稳稳扛住了根本不需要在网关层加额外逻辑。audit参数置为True后每次调用会输出一条审计日志记录调用时间、来源、方法名、参数摘要和返回状态。对排查线上问题很有帮助缺点是日志量大内部非必要场景建议手动开启。4.3 客户端调用参数客户端的RPCClient构造参数比较直观参数作用host服务端IPport服务端端口timeout单次调用的超时时间默认继承服务端不对客户端单独设retry失败重试次数retry_interval每次重试之间的间隔秒数keepalive_interval连接保活心跳间隔auth_token与服务端匹配的鉴权令牌客户端和服务端的超时是两张皮必须分别确认。客户端说“这个请求最多等3秒”服务端却说“我这个方法最长执行5秒”结果就是客户端3秒就断开了服务端还把任务跑完了响应回来无人接任务其实执行成功但客户端拿到了超时错误。这个坑我一开始就被绊倒过。5. 实战案例把内部图片压缩服务改造成a2rpc这个案例是我在生产环境做的第一次完整迁移。原来的架构是一个采集服务不断抓取图片抓完以后调用一个独立的HTTP压缩服务压缩完成后上传对象存储。瓶颈有两个HTTP接口每次请求都要带着完整的多部分表单采集服务要等压缩服务响应之后才能抓下一张。整体吞吐量上不去。5.1 改造前的调用痛点改造前的调用链大致是# 改造前伪代码 for image_path in image_list: resp requests.post( http://compress-api/internal/compress, files{file: open(image_path, rb)}, timeout30, ) result resp.json()每一张图片都要重新建立TCP连接、上传完整文件、等待响应。压缩本身只花几百毫秒但加上网络传输、协议解析、连接建立时间单张耗时拉到两秒上下。最棘手的是采集器的爬取速度和压缩服务的能力不匹配一边在快速抓取一边在排队请求。5.2 改造后的服务端设计我把压缩逻辑封装成一个CompressService类内部维护一个有限大小的线程池专门执行图片压缩这类CPU密集型操作# compress_service.py import asyncio, time, base64 from concurrent.futures import ThreadPoolExecutor from a2rpc import RPCServer, rpc_method from img_compressor import compress_image_bytes class CompressService: def __init__(self): self._executor ThreadPoolExecutor(max_workers8) rpc_method(timeout20) async def compress_png(self, image_bytes_b64: str, quality: int 80) - str: raw base64.b64decode(image_bytes_b64) loop asyncio.get_running_loop() compressed await loop.run_in_executor( self._executor, compress_image_bytes, raw, quality, ) return base64.b64encode(compressed).decode() async def main(): server RPCServer( CompressService(), bind0.0.0.0, port9101, health_check_port9102, ) await server.serve() asyncio.run(main())这个设计有几点用心图片以base64字符串形式在JSON化的消息里传递避免设计复杂的二进制分帧传输。把真正的压缩任务丢到线程池执行避免阻塞事件循环。timeout20给了压缩一定的裕量因为大图片压缩可能不止几秒。健康检查端口独立方便接入容器探活。5.3 客户端批量并发调用压缩服务上线后采集端的调用变成这样# client_batch.py import asyncio from a2rpc import RPCClient BATCH_SIZE 16 async def compress_one(client, image: bytes): import base64 payload base64.b64encode(image).decode() result_b64 await client.compress_png( image_bytes_b64payload, quality85, ) return base64.b64decode(result_b64) async def run_batch(images): async with RPCClient(127.0.0.1, 9101, timeout25, retry2) as client: semaphore asyncio.Semaphore(BATCH_SIZE) async def guarded(image): async with semaphore: return await compress_one(client, image) tasks [asyncio.create_task(guarded(img)) for img in images] return await asyncio.gather(*tasks)压测结果很清楚32张图片的批量任务改造前串行HTTP耗时约64秒改造后并发16路整体耗时约7秒吞吐量提升接近九倍。主要时间花在了压缩本身和少量网络传输上等待时间被并发吃掉了。这个案例里a2rpc的实际价值不在于它把延迟消除了而在于它让“并发调用远程服务”的复杂度降下来了。写起来四五行代码不需要自己维护连接池和请求队列。5.4 压测与调参过程我想多说一句压测里看到的真实数据。刚开始我顺手配置了RPCClient(..., timeout3)结果大量报TLE。排查发现单张高清图片压缩时间本身就超过3秒客户端超时设得太紧任务白白执行了却拿不到结果。后来把客户端超时放宽到25秒把服务端装饰器超时也统一成20秒再配合客户端的重试机制整个链路的稳定性才真正建立起来。这背后的原则是一致的客户端超时一定大于服务端最大可接受执行时间而且两者差距要留足网络传输余量。6. 排错手记五个让我挠过头皮的坑逐个还原排查链路排错章节我放在最后一部分之前是因为这些坑几乎每个迁移者都会遇到。这里我不直接给结论而是还原我自己排查的步骤方便大家参考。6.1 第一个坑客户端超时报错服务端还在默默干活现象客户端等待15秒后抛出超时错误服务端日志毫无异常业务上任务的真实结果其实是成功的。排查链路先看客户端日志确认超时时间。再看服务端方法执行日志确认函数确实被调起了。对比两端日志时间戳发现服务端完成时间比客户端超时时间晚两秒左右。查看装饰器timeout发现服务端方法本身上限是20秒客户端只给了15秒。解决调整客户端超时到30秒服务端超时设置到25秒并固定成配置项。这个坑完全是因为两端各自独立配置惹出来的我在项目里后来强制约定了一套规则所有RPC调用配置项必须以环境变量统一注入避免开发环境改一处、生产环境漏一处。6.2 第二个坑大消息体被静默截断现象服务端接收一个很大的列表参数时客户端报错说解包失败服务端毫无反应。排查链路客户端本地打印参数长度没发现异常。打开调试日志确认消息字节数约1.6MB。翻源码发现底层有消息体尺寸限制默认在1MB左右。解决在RPCServer构造时调整max_request_size参数按实际场景调大到10MB。同步把max_response_size也调了。这个坑提醒我一个原则方法参数体量大的时候先确认传输上限不要默认能传大对象。6.3 第三个坑Windows开发环境连不上服务端现象本机Windows开发环境客户端连接服务端始终失败但Linux容器内运行同一套客户端没问题。排查链路检查防火墙窗口弹窗全放行。检查服务端是不是监听在0.0.0.0确认无误。用telnet验证端口确实可连通说明网络层没问题。翻a2rpc底层实现确认它用的是asyncio自带的StreamReader/StreamWriter。想起Windows上默认事件循环策略不同手动切换事件循环策略后复测问题消失。解决在Windows端启动脚本里调用asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())。跨平台有差异这种事真的要亲历一次才记得住。6.4 第四个坑客户端开了多个连接导致文件描述符泄漏现象长跑任务运行几小时候客户端进程抛出文件描述符不足的错误。排查链路lsof -p pid | grep TCP查看客户端建立的大量TCP连接。发现代码内部循环里每次都新建了RPCClient没有复用连接。解决整个进程内只保留一个RPCClient实例用连接池复用连接只在进程退出时才关闭。这类问题在HTTP时代也会遇到但RPC因为是长驻连接问题潜伏期更短更加显眼。6.5 第五个坑多服务类注册时方法名冲突现象注册了两个服务类一个命名为UserService另一个也带个叫get_info的方法结果客户端调用A服务一直进错B服务。排查链路在服务端打印注册的方法路由表发现两个同名方法都注册了。后来者覆盖了前者的路由。解决在注册时加上prefix参数把不同模块的方法命名空间强制隔开。设置prefix为user后调用接口变成user_get_info冲突消失。这一步也是在读文档时顺手看到的等到真踩了坑才真正觉得重要。7. 除了纯RPC服务a2rpc还能这样融入现有项目跑通上面的服务之后很多人会问那我的项目已经上了FastAPI、Django或者Celery还有必要为这几个内部方法再开一个RPC端口吗我的答案取决于场景。这里提供三种可行的融合姿势都是我在现有系统里实际用过的。7.1 和FastAPI并存RPC扛批量HTTP扛面向外部API我的实践方式是让它俩并行存在FastAPI负责对外提供REST接口a2rpc只对内部节点开放。理由很朴素外部开放接口需要认证、API Key、流控这些FastAPI生态更成熟而内部服务之间高频调用的小方法用RPC更省心。两者的健康探活可以共用同一个K8s podRPC服务单独开放健康检查端口即可。两个服务在同一个进程内共存时只需要确保FastAPI和RPCServer共用一个事件循环做法如下import uvicorn from a2rpc import RPCServer from fastapi import FastAPI app FastAPI() rpc_server RPCServer(CalcService(), bind0.0.0.0, port9100) app.on_event(startup) async def start_rpc(): await rpc_server.start() app.on_event(shutdown) async def stop_rpc(): await rpc_server.stop()两个服务都跑在asyncio的同一个loop中协程间可以畅通无阻。7.2 和Celery并存异步任务里调用RPC方法如果团队里已经用了Celery做异步任务队列RPC服务作为执行节点接入也非常平滑。Celery的Worker进程启动时直接把RPC服务挂在后台即可。# celery_app.py from celery import Celery from a2rpc import RPCServer from compress_service import CompressService celery_app Celery(tasks, brokerredis://...) rpc_server RPCServer(CompressService(), bind0.0.0.0, port9101) celery_app.task def start_rpc_worker(): # 调用RPC服务使Celery Worker同时充当RPC服务端 asyncio.run(rpc_server.serve())但这种方式需要注意进程竞争Celery Worker默认fork策略和asyncio的loop创建时机比较敏感不要在有共享连接池的前提下再fork。生产环境我会让Celery worker单独一个进程RPC服务的生命周期独立管理避免耦合。7.3 单机多进程架构里做IPC还有一个被忽略的使用场景单机多进程服务之间互相协调时RPC也可以作为IPC替代方案。比如爬虫采集节点、AI推理worker、或数据导出进程都在同一台宿主机运行它们之间的通信可以用named pipe或Unix socket但用a2rpc的好处是将来扩展到多机时代码几乎不用改只需要把客户端地址从本机IP改成远程IP。# 本机IPC示例 rpc_server RPCServer(TaskService(), bind127.0.0.1, port9300) client RPCClient(127.0.0.1, 9300)从这个角度看a2rpc提供的是一个通信协议与编解码层而不是强绑定网络拓扑。8. 我现在的选择标准和最后一点建议经历了多次迁移后我对一个内部服务到底需不需要上RPC有了比较清晰的选择标准如果只是给外部客户端提供几个接口又需要API文档、鉴权和限流继续用HTTP框架没必要引入RPC。如果是在内部多个Python进程/节点间高频调用方法且调用方和被调用方都以Python为主a2rpc是非常合适的低成本方案。如果对性能有极致要求或者做成了跨语言的对外标准化服务RPC会显得不够用这时就得考虑gRPC这类完整框架。关于选型对比我用一张表总结了实际感受方案序列化代码量跨语言内部Python进程间效率原始asyncio socket自定很大可以但繁琐一般HTTP JSONJSON中好低gRPCProtobuf很大好中需要代码生成a2rpcmsgpack/JSON极小受限高这里我并不是说a2rpc多完美更准确的说法是在“全Python环境、异步、内部服务”这个特定场景里它击中了效率与代码量的平衡点。最后分享一点实际经验不管用什么RPC库第一件事不是赶着写功能代码而是先列一张“可调用方法清单”写明方法名、参数、返回值、超时时间、异常语义。跑偏的RPC迁移多半是从方法设计师期四处漏风、后期排错众里寻他开始的。你自己去试的时候建议先搭一个最小服务端起一个最小客户端把5.3节的批量并发模式复制过去先稳定跑通一批小任务再逐步换大参数、加大数据一切性能上的问题都会在你面前慢慢显出原形。