
做后端时间长了你会发现一个很现实的问题业务量一上来单体应用必然面临拆分。服务一拆通信就成了绕不开的事。服务A要调服务B的数据总不能直接连对方的数据库更不可能把代码复制一份贴过来。这时候你得有个像样的远程调用方案于是在Python后端开发里RPC就成了必修课。这一篇是“Python后端开发之旅”系列的第三篇我打算把两件事一次讲透服务间通信用的RPC以及处理耗时任务的异步方案。这两样东西在后端架构里经常成对出现一个负责“把请求送出去”一个负责“把重活放后台干”。我会结合自己的实战经验从概念讲到代码实现最后再整理一些你迟早会遇到的坑。适合谁看两类人。一类是刚接触微服务、搞不清RPC和普通HTTP接口区别的Python后端新人另一类是已经写了不少接口但对服务治理和异步任务总是似懂非懂、希望系统梳理一遍的开发者。不管你是哪一类这篇文章我尽量少讲空话直接上可落地的方案。1. 为什么后端服务都绕不开RPC1.1 从单体拆分开始说起先说一个我印象很深的事。刚工作那会儿公司核心系统是个巨大的单体应用所有功能都在同一个进程里跑。模块之间调用很简单直接import函数就行效率高也没啥网络开销。但后来业务膨胀几十人维护同一套代码每次上线都提心吊胆改一个模块可能把整个服务搞崩。于是架构师决定拆分服务。服务一拆第一个问题就来了原本的函数调用变成了服务间调用。比如订单服务要查用户信息用户服务在另一台机器上代码没法直接import只能通过网络请求。这就要引入远程调用方案。在当时最常见的选择就是HTTP接口但很快发现接口多了以后性能问题突出而且服务间调用的语义对不上。我说不好听的HTTP接口是给人看的RPC天生就是给程序用的。在这个背景下RPCRemote Procedure Call远程过程调用才是更贴近函数调用体验的方案。它的核心思想很简单让调用方像调用本地函数一样调用远程服务参数和返回值在进程间传输序列化和网络通信的细节全部封装掉。1.2 RPC与HTTP到底差在哪很多人会问我直接用HTTP JSON不也挺好的省事到处都能用。这话对一半。对外提供的OpenAPI确实适合HTTP因为要兼容各种语言的客户端通用性第一。但服务内部之间的调用场景不一样追求的是低延迟、高吞吐、强类型约束这时候HTTP JSON的短板就暴露了。我说三个最直接的区别。第一传输效率和性能。HTTP/1.1 JSON是文本协议头部信息冗长JSON序列化也有明显开销。RPC框架通常用二进制序列化最典型的是Protobuf序列化后的体积可能只有JSON的十分之一左右解析速度也快得多。在高并发内部调用场景下这个差距会被无限放大。第二强类型契约。HTTP JSON是弱类型的字段类型对不对得靠人约定改接口文档很容易踩坑。RPC框架通常基于IDL接口描述语言定义服务接口、请求参数和响应结构编译时就能生成代码类型不匹配直接在编译阶段暴露前后端联调就像写同一个项目。第三服务治理能力。RPC框架普遍内置负载均衡、超时控制、熔断、重试、链路追踪这些能力而HTTP接口要自己造轮子或者额外接一套网关。我做了一个对比表方便你看关键差异维度HTTP JSONRPC以gRPC为例数据格式文本JSON二进制Protobuf传输协议HTTP/1.1/2HTTP/2接口约束弱类型文档约定强类型IDL定义序列化性能较慢很快服务治理依赖外部组件框架内置排错友好度直观易读需要解码工具适用场景对外API、跨语言泛化调用内部服务间高性能调用1.3 常见的Python RPC框架怎么选Python生态里RPC框架其实不少但主流的是这五个gRPC、Thrift、JSON-RPC、XML-RPC、以及一些公司自研的框架。我实际用下来最推荐的是gRPC其次看场景选Thrift。gRPC是Google开源的基于HTTP/2默认用Protobuf序列化。它的优势非常明显生态成熟跨语言支持强大Python客户端和服务端实现都很稳官方文档齐全社区案例多。我这边多个服务间通信都是gRPC压测数据也很满意单个请求的平均耗时比HTTP方案低了一半还多。Thrift是Facebook开源的性能也不错但Python的社区活跃度和工具链明显不如gRPC如果你不是有特殊的历史包袱我不太建议新项目从gRPC跨到Thrift。JSON-RPC是基于JSON的轻量级RPC协议优点就是简单调试直观适合内部小规模工具类调用但复杂场景下服务治理能力偏弱。还有一组选择是直接用消息队列做请求响应模式比如RabbitMQ RPC模式这其实是塞进了异步消息的范畴后面说异步任务的时候会提到。我这边的经验法则就一条内部服务间通信优先gRPC对外API用HTTP JSON需要广播、事件驱动用消息队列。别想着一个方案打通所有场景。2. 用gRPC打通服务间通信2.1 proto文件先定义契约再写代码gRPC的第一步是写proto文件这文件就是服务之间的“合同”。写好了proto前后端、多语言客户端就能基于同一份契约各自生成代码从源头上避免沟通不一致。我以视频处理服务为例定义一个简单的服务接口。假设我们有一个视频转码服务客户端传入视频信息服务端返回任务状态。syntax proto3; package video; // 定义视频处理任务的数据结构 message VideoRequest { string video_id 1; string input_path 2; string output_path 3; int32 width 4; int32 height 5; } message TaskStatus { string task_id 1; string state 2; // PENDING / STARTED / SUCCESS / FAILURE int32 progress 3; // 0 - 100 } // 定义视频处理服务 service VideoProcessor { rpc SubmitVideo(VideoRequest) returns (TaskStatus); rpc GetTaskStatus(TaskStatusRequest) returns (TaskStatus); } message TaskStatusRequest { string task_id 1; }写proto有几个注意点字段编号不能随便改proto3序列化依赖字段编号一旦接口上线字段编号变了会导致旧客户端解析失败。我见过有人中途改字段编号结果线上数据乱套的事故。字段类型要用proto自带类型int32、int64、string、bytes、repeated这些不要直接用Python的字典等类型去硬套。rpc方法设计要考虑调用模式gRPC支持四种简单RPC请求-响应、服务端流式、客户端流式、双向流式。大部分内部调用用简单RPC就够了。2.2 生成代码与服务端实现proto文件写完之后要用工具生成Python代码。安装依赖和生成命令是这样的pip install grpcio grpcio-tools protobuf # 在proto文件所在目录执行 python -m grpc_tools.protoc -I. --python_out. --grpc_python_out. video.proto执行完以后目录下会多出两个文件video_pb2.py消息类型的序列化代码和video_pb2_grpc.py服务端和客户端的桩代码。这两个文件都不建议手改生成后直接引用就行。接下来实现服务端逻辑。继承生成的VideoProcessorServicer重写对应的rpc方法import grpc from concurrent import futures import video_pb2 import video_pb2_grpc class VideoProcessorService(video_pb2_grpc.VideoProcessorServicer): def SubmitVideo(self, request, context): # 这里是处理视频提交的逻辑 # 注意真实场景不要在这里做耗时操作后面讲异步任务时会说明 print(f收到视频任务: {request.video_id}) return video_pb2.TaskStatus( task_idtask_001, statePENDING, progress0 ) def GetTaskStatus(self, request, context): # 查询任务状态通常从数据库或缓存读取 return video_pb2.TaskStatus( task_idrequest.task_id, stateSUCCESS, progress100 ) def serve(): server grpc.server(futures.ThreadPoolExecutor(max_workers10)) video_pb2_grpc.add_VideoProcessorServicer_to_server( VideoProcessorService(), server ) server.add_insecure_port([::]:50051) print(gRPC server started on port 50051) server.start() server.wait_for_termination() if __name__ __main__: serve()这里有个细节要注意ThreadPoolExecutor的max_workers要结合你的服务类型设置。如果是CPU密集型服务线程多了反而容易触发GIL竞争如果是IO密集型比如大量数据库查询线程数可以适当调大。我常用的经验是IO密集型max_workers在CPU核数的2倍到4倍之间调优CPU密集型推荐用进程池方案或者直接上gRPC的异步接口。2.3 客户端调用与连接管理gRPC客户端的调用方式非常简洁但连接管理上有个坑很多人会踩。import grpc import video_pb2 import video_pb2_grpc channel grpc.insecure_channel(localhost:50051) stub video_pb2_grpc.VideoProcessorStub(channel) # 提交视频任务 response stub.SubmitVideo( video_pb2.VideoRequest( video_idvideo_12345, input_path/data/input/a.mp4, output_path/data/output/a_720p.mp4, width1280, height720 ) ) print(response.task_id, response.state, response.progress)坑就在这里很多人每次调用都新建一个channel用完也不关。gRPC的channel建立需要握手协商成本很高高并发场景下会拖慢整体响应。正确做法是复用channel。我自己处理的方式是建一个全局连接池或者用单例模式持有一个channel设置合理的keepalive参数。另外不要忘了超时控制gRPC默认的deadline是None也就是不设超时一旦远程服务挂掉调用方可能一直阻塞。建议调用时都加上timeout参数response stub.SubmitVideo( video_pb2.VideoRequest(video_idvideo_12345), timeout10 # 单位秒 )这样远程服务无响应时10秒后客户端就会收到一个超时异常不至于把调用方的线程池拖垮。3. 异步任务把耗时操作从RPC里搬出去3.1 同步调用为什么会卡死先讲一个反面教材。早期我做一个视频处理服务gRPC的SubmitVideo接口直接同步执行转码逻辑。转码是个重型操作一个1GB的视频转码加压缩快则二三十秒慢的话可能几分钟。第一次压测就出事了客户端调用SubmitVideo后一直没响应最后报了个错——cannot finish rpc call in 30 seconds。原因是客户端设置了10秒超时服务端处理迟迟不返回30秒后中断了调用。问题很明显耗时任务不能放在RPC的同步处理函数里。RPC适合做短平快的请求响应比如提交任务、查询状态、查询配置。真正的重活应该用异步任务队列把任务投递到后台让worker慢慢处理。这也是我在开头说的RPC和异步任务经常成对出现的原因。一个典型的模式是RPC接口快速接收请求把任务入队立刻返回一个任务ID客户端拿着这个ID去轮询状态或者等服务端回调通知。3.2 Celery核心概念与配置Python生态里最成熟的异步任务框架是Celery老牌、稳定、功能全面支持多队列、定时任务、任务重试、结果存储。相比轻量级的RQ和DramatiqCelery的学习曲线稍微陡一点但功能上限高团队协作时通用性也更好。Celery有四个核心概念Broker是消息中间件负责接收任务消息并派发给worker。生产环境最常见的选择是Redis或RabbitMQ。Redis配置简单、性能好适合中小规模项目RabbitMQ稳支持更多消息特性适合对可靠性要求更高的场景。Backend是结果存储保存任务的执行结果方便查询状态。如果不需要保存任务结果可以不用配backend只做Fire-and-forget。但大多数场景下我们至少要记录任务状态所以backend是标配。Worker是真正干活的工作进程从broker里拉取任务消息并执行。Queue是任务队列一个worker可以监听多个队列不同队列可以分配给不同优先级的worker。我自己常用的配置如下# celery_app.py from celery import Celery app Celery( video_tasks, brokerredis://localhost:6379/0, backendredis://localhost:6379/1 ) app.conf.update( # 指定任务序列化方式为JSON避免使用默认的pickle task_serializerjson, accept_content[json], result_serializerjson, timezoneAsia/Shanghai, enable_utcTrue, # 单个任务硬超时单位秒 task_time_limit600, # 单个任务软超时超时后会抛异常让代码自行处理 task_soft_time_limit540, # 不显示任务执行日志可选 worker_hijack_root_loggerFalse, )注意task_serializer我强制用JSON。为什么不推荐默认的pickle因为pickle在Python之间传递对象很方便但存在安全性问题——如果broker被攻破恶意数据会被反序列化执行危险操作。另外pickle序列化后的消息体积通常比JSON大网络传输效率也低。3.3 任务定义与Worker启动定义任务用app.task装饰器。一个实际的视频转码任务大概是这样的# tasks.py from celery_app import app import subprocess app.task(bindTrue, namevideo.transcode, max_retries3, default_retry_delay60) def transcode_video(self, video_id, input_path, output_path, width, height): try: # 模拟耗时的转码逻辑实际场景用ffmpeg命令处理 cmd [ ffmpeg, -i, input_path, -vf, fscale{width}:{height}, output_path ] subprocess.run(cmd, checkTrue, capture_outputTrue) return { video_id: video_id, status: completed, output_path: output_path } except subprocess.CalledProcessError as exc: # 失败自动重试最多3次每次间隔60秒 raise self.retry(excexc)这里我用bindTrue把任务实例绑定到第一个参数self这样可以在任务内部调用self.retry实现重试。重试参数max_retries和default_retry_delay控制重试次数和间隔。我的经验是网络波动、临时性资源不足导致的失败适合重试但业务参数本身错误导致的失败重试没有意义反而会占用队列资源这种情况要在任务开头做好参数校验直接抛出不重试的异常。任务写好后启动worker的命令很简单celery -A tasks worker --loglevelinfo -c 4 -Q default-c指定并发进程数一般建议和CPU核心数对齐CPU密集型任务设成核数即可IO密集型任务可以适当多开。我这边一般先跑默认值看监控再调。调用方式就更简单了from tasks import transcode_video result transcode_video.delay( video_12345, /data/input/a.mp4, /data/output/a_720p.mp4, 1280, 720 ) # result.id 就是任务ID后续查状态用用delay方法直接把任务投递到broker函数立刻返回后台worker异步执行。任务ID可以存到数据库里前端轮询状态时使用。4. 一个完整的RPC异步任务实战4.1 场景设计搭建一个视频处理服务前面讲了RPC和Celery各自的基本用法这一节把它们串起来搭建一个完整的视频处理服务。我们的目标场景是客户端上传视频文件服务端返回任务ID客户端通过任务ID查询转码进度。这个场景非常典型很多文件处理类服务都是这个套路。整体处理流程是客户端调用SubmitVideo RPC接口传入视频ID和转码参数gRPC服务端收到请求后先把任务写入数据库状态PENDING然后通过celery把转码任务投递到brokergRPC接口立即返回任务ID给客户端worker从broker拉取任务执行实际转码执行过程中更新数据库中的进度和状态客户端定期调用GetTaskStatus RPC接口查询状态。这样的架构里有一个关键设计思想RPC接口和worker是解耦的。RPC接口只负责接收请求、快速响应具体业务逻辑全部丢给异步任务互不影响。即使worker挂了任务还在broker里排队等worker恢复后继续执行不会因为进程重启而丢任务。4.2 从RPC接收请求到投递任务gRPC服务端这一步的代码需要改一下把原来同步转码变成任务投递。# video_processor.py import grpc from concurrent import futures import video_pb2 import video_pb2_grpc from tasks import transcode_video from database import create_task_record, update_task_status class VideoProcessorService(video_pb2_grpc.VideoProcessorServicer): def SubmitVideo(self, request, context): # 1. 写入任务记录到数据库初始状态PENDING task_id create_task_record( video_idrequest.video_id, input_pathrequest.input_path, output_pathrequest.output_path, widthrequest.width, heightrequest.height, statePENDING ) # 2. 投递异步任务参数里带上task_id transcode_video.delay( task_id, request.video_id, request.input_path, request.output_path, request.width, request.height ) # 3. 立即返回任务ID和状态 return video_pb2.TaskStatus( task_idtask_id, statePENDING, progress0 ) def GetTaskStatus(self, request, context): # 从数据库读取任务信息 task_info get_task_record(request.task_id) return video_pb2.TaskStatus( task_idrequest.task_id, statetask_info[state], progresstask_info[progress] )这里每一步都很关键。数据库记录是任务状态的中心RPC接口快速写入和读取celery worker在后台更新。三个模块共享同一个数据库状态保持一致。worker里的任务也要改把结果写入数据库# tasks.py from celery_app import app from database import update_task_status app.task(bindTrue, namevideo.transcode, max_retries3, default_retry_delay60) def transcode_video(self, task_id, video_id, input_path, output_path, width, height): try: update_task_status(task_id, stateSTARTED, progress10) # 实际转码命令这里用subprocess调ffmpeg run_ffmpeg_transcode(input_path, output_path, width, height) update_task_status(task_id, stateSUCCESS, progress100) return {task_id: task_id, status: completed} except TranscodeError as exc: update_task_status(task_id, stateFAILURE, progress-1) raise self.retry(excexc)更新进度这种细节实际项目里可以在ffmpeg执行过程中解析进度输出定期更新到数据库。这样做的好处是前端能看到具体的进度百分比体验会好很多。4.3 任务状态查询与回调设计任务状态查询的逻辑上面已经写了就是GetTaskStatus这个RPC方法。但除了轮询还有一种更主动的通知方式回调。回调有两种常见实现。一是任务完成后worker通过HTTP POST或gRPC调用通知客户端指定的回调地址。缺点是需要客户端提供一个可访问的接口在公网环境下还有安全认证的问题。二是通过WebSocket推送状态。服务端维持WebSocket连接任务状态变化时主动推送给前端。这种方式实时性好但实现复杂度高需要在服务端维护连接状态。我的经验是内部系统优先用轮询简单可靠成本低对外提供的服务如果一定要实时通知再考虑WebSocket方案。轮询的频率也不要太高2到5秒一次就足够了频率太高会给服务端造成不必要的压力。5. 实战中躲不开的坑5.1 gRPC超时问题实录前面提到的cannot finish rpc call in 30 seconds这个报错我后来做了详细排查。原因不只是客户端超时还有一个隐藏问题gRPC服务端如果长时间不返回客户端默认有个deadline机制隔一段时间会主动断开。即便客户端没有显式设置timeoutgRPC框架在特定情况下也有默认限制。遇到这类问题先别急着调大timeout参数。根本解法是重构把耗时操作移出RPC处理流程换成异步任务模式。RPC接口只做两件事——接收参数、投递任务然后立刻返回。这样RPC方法执行时间往往在几十毫秒内远不会触发超时。如果你确认某个RPC方法确实要执行长时间操作比如批处理任务那就把deadline显式设置大一点比如300秒并且要做好日志埋点方便追踪到底慢在哪里。5.2 Celery常见问题速查Celery用久了有些问题几乎每换一个项目就会碰到。我整理了最常见的几个报错或现象可能原因解决方案任务一直不执行队列堆积broker配置错误或worker未启动检查broker地址确认worker命令执行无误任务执行两遍没有开启ack早确认机制任务超时后被重新投递设置task_acks_late和task_reject_on_worker_lost同时任务做好幂等导入错误Could not import moduletasks.py依赖了不在Python路径上的包确保启动worker时在项目目录下或设置PYTHONPATH序列化错误TypeError任务参数传了自定义对象把所有任务参数改成基本类型或字典worker内存涨得很快存在内存泄漏或并发过高排查代码降低并发数使用prefork模式别用threads模式第一列的“任务执行两遍”是个大坑。默认情况下Celery在worker收到任务后就发送ack确认如果任务执行中途worker崩溃broker会重新投递这个任务造成重复执行。解决方案是把acks设置为late也就是任务执行完成后再确认。但这带来另一个要求任务必须幂等重复执行的结果要一样否则还是会有问题。做任务设计时优先保证幂等性。5.3 性能优化经验分享最后聊几个我实测有效的优化经验。第一个是gRPC的channel复用前面说过。你一定不要每次请求新建channel这个优化对整体性能提升非常明显。压测中复用channel的QPS比每次新建channel高出好几倍。第二个是Celery的prefetch_multiplier参数。默认值是4意思是每个worker进程每次可以预取4个任务。如果任务执行时间比较短预取可以提升吞吐但如果任务执行时间很长预取太多会导致任务分配不均衡——一台worker积压任务另一台闲着。处理耗时较长的任务时我会把这个值调成1celery -A tasks worker --loglevelinfo -c 4 -Q default --prefetch-multiplier1第三个是队列分拆。不要把所有的任务都丢到同一个队列里要按重要程度和资源需求分队列。比如转码任务放到transcode队列发通知的轻量任务放到default队列然后启动不同数量的worker分别消费。这样即使转码任务堆积也不会影响消息通知的及时性。第四个是任务结果的存储策略。很多人习惯把任务结果全都存到Redis时间长了Redis内存膨胀。我一般设置result_expires参数比如保存2小时就够了app.conf.result_expires 7200 # 单位秒超过2小时的结果自动清理避免Redis变成垃圾场。回到开头那个项目这套RPC异步任务的架构上线后接口超时问题基本绝迹视频转码吞吐量也翻了倍。后来我总结后端开发里很多看起来很难的问题本质上都是“同步阻塞”造成的。RPC解决的是远程调用的问题异步任务解决的是耗时操作的问题两者配合服务的弹性和稳定性都会上一个台阶。从个人体会来说学RPC和异步任务一定要动手做一个完整的项目不要只跑demo。你真正理解了“RPC快速接收 异步慢慢处理 状态回查”这种组合模式后端服务设计就通了一大半。后面的文章我打算再深入讲讲消息队列的可靠投递和服务治理的话题咱们下篇见。