如何从Python客户端向C# gRPC服务端发送CancellationToken取消请求
问题描述
我本以为这是个容易解决的小问题,但实际并非如此。我正在尝试实现从Python客户端向C# gRPC服务端发送CancellationToken取消信号。
在C#客户端中,我可以这样实现:
var tokenSource = new CancellationTokenSource(); while (await priceResponse.ResponseStream.MoveNext(tokenSource.Token))
服务端代码如下:
public async override Task Prices(PricesRequest request, IServerStreamWriter<PricesResponse> responseStream, ServerCallContext context) { ... while (!context.CancellationToken.IsCancellationRequested) { ... }
我曾尝试自定义CancellationToken类,但没有效果:
class CancellationToken: def __init__(self): self._is_canceled = False def is_canceled(self): return self._is_canceled def cancel(self): self._is_canceled = True
以下是我的Python客户端代码:
import grpc import greeting_pb2 import greeting_pb2_grpc import pricing_pb2 import pricing_pb2_grpc import os import logging _LOGGER = logging.getLogger(__name__) _SERVER_HOST = "localhost" _SERVER_PORT = "50051" def pricing_client(): try: _LOGGER.info(f"Attempting to connecto {_SERVER_HOST}:{_SERVER_PORT}...") channel = grpc.insecure_channel(f"{_SERVER_HOST}:{_SERVER_PORT}") pricingStub = pricing_pb2_grpc.PricingServiceStub(channel) _LOGGER.info("Channel Connection made, stubs created...") myrequest = greeting_pb2.GreetingRequest(greeting=mygreeting) myprice = pricing_pb2.Price(instrument="BP") mypricerequest = pricing_pb2.PricesRequest(price=myprice) pricingResponse = pricingStub.Prices(mypricerequest) _LOGGER.debug(f"Pricing response executed successfully {pricingResponse}") prices = 0 for price in pricingResponse: _LOGGER.debug(f"Price received: {price}") prices += 1 if(prices == 10): # trying to cancel here, but for now I just exit the client return except: _LOGGER.error(f"An error occured - please check the server is runing @ {_SERVER_HOST}:{_SERVER_PORT}") finally: _LOGGER.info("Application exiting...") if __name__ == "__main__": logging.basicConfig(level=logging.DEBUG) pricing_client()
请问如何用Python实现该取消模式与C#服务端通信?
解决方案
自定义的CancellationToken类仅在客户端本地生效,无法将取消信号传递给gRPC服务端,需要使用gRPC Python原生的取消机制。对于服务器流式调用,可通过with_call()方法获取gRPC调用对象,调用其cancel()方法发送取消信号,该信号会被C#服务端的context.CancellationToken捕获。
修改后的Python客户端代码:
import grpc import greeting_pb2 import greeting_pb2_grpc import pricing_pb2 import pricing_pb2_grpc import os import logging _LOGGER = logging.getLogger(__name__) _SERVER_HOST = "localhost" _SERVER_PORT = "50051" def pricing_client(): try: _LOGGER.info(f"Attempting to connect to {_SERVER_HOST}:{_SERVER_PORT}...") channel = grpc.insecure_channel(f"{_SERVER_HOST}:{_SERVER_PORT}") pricingStub = pricing_pb2_grpc.PricingServiceStub(channel) _LOGGER.info("Channel Connection made, stubs created...") mygreeting = "test" myrequest = greeting_pb2.GreetingRequest(greeting=mygreeting) myprice = pricing_pb2.Price(instrument="BP") mypricerequest = pricing_pb2.PricesRequest(price=myprice) # 使用with_call()获取响应迭代器和gRPC调用对象 pricingResponse, call = pricingStub.Prices.with_call(mypricerequest) _LOGGER.debug(f"Pricing response executed successfully {pricingResponse}") prices = 0 for price in pricingResponse: _LOGGER.debug(f"Price received: {price}") prices += 1 if prices == 10: # 发送取消信号到服务端 call.cancel() _LOGGER.info("Cancel signal sent to server") return except grpc.RpcError as e: if e.code() == grpc.StatusCode.CANCELLED: _LOGGER.info("Call cancelled") else: _LOGGER.error(f"RPC error occurred: {e.details()}", exc_info=True) except Exception as e: _LOGGER.error(f"An error occurred - please check the server is running @ {_SERVER_HOST}:{_SERVER_PORT}", exc_info=True) finally: _LOGGER.info("Application exiting...") if __name__ == "__main__": logging.basicConfig(level=logging.DEBUG) pricing_client()
关键说明
with_call()方法会返回两个值:响应的迭代器,以及代表当前gRPC调用的grpc.Call对象- 调用
call.cancel()后,gRPC会向服务端发送取消请求,C#服务端的context.CancellationToken.IsCancellationRequested会变为true,从而终止循环 - 客户端捕获
grpc.RpcError可以针对性处理取消相关的异常,避免无差别捕获所有异常导致问题排查困难
内容的提问来源于stack exchange,提问作者PureAlpha
相关产品推荐
相关产品推荐

