这篇播客基于 Python gRPC 实现的客户端与服务端通信:
- 首先模拟一个简单的数据服务系统,客户端通过登录接口进行身份验证;在服务端验证用户名和密码后生成 Token
- 服务端接收到请求后,会通过gRPC拦截器统一完成日志记录和Token鉴权。鉴权成功后,服务端根据客户端传入的开始日期、结束日期和字段列表,使用pandas动态生成一份随机DataFrame数据,并将数据返回给客户端。
服务端代码
这里服务端注册两个主要的函数:
auth.Auth服务的login函数,用于做权限校验data.Data服务的GetData函数,用于获取数据,在这段代码中是基于用户传入的参数生成一个模拟的数据
importjsonimportloggingimportsecretsfromconcurrentimportfuturesimportgrpcimportpandasaspd logging.basicConfig(level=logging.INFO,format="%(asctime)s [%(levelname)s] %(message)s",)logger=logging.getLogger(__name__)TOKENS={}deflogin(request,_context):"""登录成功后生成并保存 Token。"""username=request.get("username")password=request.get("password")logger.info("收到登录请求:username=%s",username)ifusername!="admin"orpassword!="123456":logger.warning("登录失败:username=%s",username)return{"success":False,"message":"用户名或密码错误"}token=secrets.token_hex(16)TOKENS[token]=username logger.info("登录成功:username=%s",username)return{"success":True,"token":token,"message":"登录成功"}defget_data(request,_context):"""根据日期范围和字段生成一份随机 DataFrame 数据。"""start_date=request.get("start_date")end_date=request.get("end_date")fields=request.get("fields",[])ifnotstart_dateornotend_dateornotfields:return{"success":False,"message":"start_date、end_date、fields 不能为空"}try:dates=pd.date_range(start=start_date,end=end_date,freq="D")exceptValueError:return{"success":False,"message":"日期格式应为 YYYY-MM-DD"}dataframe=pd.DataFrame({"date":dates.strftime("%Y-%m-%d"),**{field:pd.Series(secrets.randbelow(10000)/100for_indates)forfieldinfields},})return{"success":True,"data":dataframe.to_dict(orient="records"),"message":"数据获取成功",}defjson_handler(function):"""统一处理 Python 字典和 JSON 字节之间的转换。"""returngrpc.unary_unary_rpc_method_handler(function,request_deserializer=lambdadata:json.loads(data.decode("utf-8")),response_serializer=lambdadata:json.dumps(data,ensure_ascii=False).encode("utf-8"),)classLoggingInterceptor(grpc.ServerInterceptor):"""记录每次 gRPC 请求。"""defintercept_service(self,continuation,handler_call_details):method=handler_call_details.method logger.info("收到 gRPC 请求:method=%s",method)handler=continuation(handler_call_details)ifhandlerisNone:logger.warning("找不到 gRPC 方法:method=%s",method)returnhandlerclassAuthInterceptor(grpc.ServerInterceptor):"""在业务函数执行前检查 Token。"""defintercept_service(self,continuation,handler_call_details):method=handler_call_details.method handler=continuation(handler_call_details)ifhandlerisNone:returnNoneifmethod=="/auth.Auth/Login":returnhandlerdefcheck_token(request,context):token=request.get("token")iftokennotinTOKENS:logger.warning("鉴权失败:method=%s",method)context.abort(grpc.StatusCode.UNAUTHENTICATED,"请先登录")logger.info("鉴权成功:method=%s, user=%s",method,TOKENS[token])returnhandler.unary_unary(request,context)returnhandler._replace(unary_unary=check_token)defcreate_server(address="[::]:50051",max_workers=4):"""创建并配置 gRPC 服务。"""grpc_server=grpc.server(futures.ThreadPoolExecutor(max_workers=max_workers),interceptors=[LoggingInterceptor(),AuthInterceptor()],)grpc_server.add_generic_rpc_handlers((grpc.method_handlers_generic_handler("auth.Auth",{"Login":json_handler(login)},),grpc.method_handlers_generic_handler("data.Data",{"GetData":json_handler(get_data)},),))grpc_server.add_insecure_port(address)returngrpc_serverdefserve():"""启动 gRPC 服务并持续等待请求。"""address="[::]:50051"grpc_server=create_server(address)grpc_server.start()logger.info("gRPC 服务已启动:address=%s",address)grpc_server.wait_for_termination()if__name__=="__main__":serve()客户端代码
客户端是极简代码,直接登录,获取token后直接获取数据即可
importjsonimportgrpc# 把 Python 字典转换成 JSON 发送,并把 JSON 响应转换回来。defcall(channel,method,data):request=channel.unary_unary(method,request_serializer=lambdavalue:json.dumps(value).encode("utf-8"),response_deserializer=lambdavalue:json.loads(value.decode("utf-8")),)returnrequest(data)# 连接服务端。withgrpc.insecure_channel("127.0.0.1:50051")aschannel:# 1. 登录,获取 Token。login_result=call(channel,"/auth.Auth/Login",{"username":"admin","password":"123456"},)print(login_result["message"])# 2. 使用 Token 查询指定日期和字段的数据。result=call(channel,"/data.Data/GetData",{"token":login_result["token"],"start_date":"2025-01-01","end_date":"2025-01-05","fields":["open","close","volume"],},)# 3. 打印服务端返回的数据。print(result["data"])