本文目录导读:

我来详细介绍Python gRPC客户端实现的完整步骤。
环境准备
首先安装必要的包:
pip install grpcio grpcio-tools protobuf
定义Proto文件
创建 helloworld.proto:
syntax = "proto3";
package helloworld;
// 服务定义
service Greeter {
// 一元RPC
rpc SayHello (HelloRequest) returns (HelloReply) {}
// 服务端流式RPC
rpc SayHelloServerStream (HelloRequest) returns (stream HelloReply) {}
// 客户端流式RPC
rpc SayHelloClientStream (stream HelloRequest) returns (HelloReply) {}
// 双向流式RPC
rpc SayHelloBidiStream (stream HelloRequest) returns (stream HelloReply) {}
}
// 请求消息
message HelloRequest {
string name = 1;
int32 age = 2;
repeated string hobbies = 3;
}
// 响应消息
message HelloReply {
string message = 1;
}
生成Python代码
python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. helloworld.proto
这会生成:
helloworld_pb2.py- 消息类helloworld_pb2_grpc.py- 服务类
基础客户端实现
import grpc
import helloworld_pb2
import helloworld_pb2_grpc
class GRPCClient:
def __init__(self, host='localhost', port=50051):
# 创建连接
self.channel = grpc.insecure_channel(f'{host}:{port}')
# 创建stub
self.stub = helloworld_pb2_grpc.GreeterStub(self.channel)
def say_hello(self, name, age=25, hobbies=None):
"""一元RPC调用"""
if hobbies is None:
hobbies = ['reading', 'coding']
# 创建请求
request = helloworld_pb2.HelloRequest(
name=name,
age=age,
hobbies=hobbies
)
try:
# 调用远程方法
response = self.stub.SayHello(request)
return response.message
except grpc.RpcError as e:
print(f"RPC failed: {e.code()} - {e.details()}")
return None
def say_hello_server_stream(self, name):
"""服务端流式RPC调用"""
request = helloworld_pb2.HelloRequest(name=name)
try:
# 获取响应流
responses = self.stub.SayHelloServerStream(request)
for response in responses:
print(f"Received: {response.message}")
yield response.message
except grpc.RpcError as e:
print(f"Stream RPC failed: {e}")
def say_hello_client_stream(self, names):
"""客户端流式RPC调用"""
def generate_requests():
for name in names:
yield helloworld_pb2.HelloRequest(name=name)
try:
response = self.stub.SayHelloClientStream(generate_requests())
return response.message
except grpc.RpcError as e:
print(f"Client stream RPC failed: {e}")
return None
def say_hello_bidi_stream(self, names):
"""双向流式RPC调用"""
def generate_requests():
for name in names:
yield helloworld_pb2.HelloRequest(name=name)
try:
responses = self.stub.SayHelloBidiStream(generate_requests())
for response in responses:
print(f"Received: {response.message}")
except grpc.RpcError as e:
print(f"Bidi stream RPC failed: {e}")
def close(self):
"""关闭连接"""
self.channel.close()
# 使用示例
def main():
client = GRPCClient('localhost', 50051)
try:
# 1. 一元RPC
result = client.say_hello('Alice', 30, ['swimming', 'reading'])
print(f"Response: {result}")
# 2. 服务端流式RPC
print("\nServer streaming:")
for msg in client.say_hello_server_stream('Bob'):
print(f" Stream received: {msg}")
# 3. 客户端流式RPC
print("\nClient streaming:")
result = client.say_hello_client_stream(['Charlie', 'David', 'Eve'])
print(f" Result: {result}")
# 4. 双向流式RPC
print("\nBidirectional streaming:")
client.say_hello_bidi_stream(['Frank', 'Grace', 'Henry', 'Ivy'])
finally:
client.close()
if __name__ == '__main__':
main()
高级客户端实现
import grpc
from typing import Optional, Callable
import helloworld_pb2
import helloworld_pb2_grpc
class AdvancedGRPCClient:
def __init__(self,
host='localhost',
port=50051,
use_ssl=False,
ssl_cert=None,
max_retries=3,
timeout=10):
self.host = host
self.port = port
self.use_ssl = use_ssl
self.max_retries = max_retries
self.timeout = timeout
self.channel = self._create_channel(ssl_cert)
self.stub = helloworld_pb2_grpc.GreeterStub(self.channel)
def _create_channel(self, ssl_cert=None):
"""创建gRPC通道"""
if self.use_ssl:
# SSL/TLS连接
with open(ssl_cert, 'rb') as f:
credentials = grpc.ssl_channel_credentials(f.read())
return grpc.secure_channel(f'{self.host}:{self.port}', credentials)
else:
# 不安全连接
return grpc.insecure_channel(f'{self.host}:{self.port}')
def _with_retry(self, func, *args, **kwargs):
"""带重试机制的调用"""
for attempt in range(self.max_retries):
try:
return func(*args, **kwargs)
except grpc.RpcError as e:
if attempt == self.max_retries - 1:
raise
print(f"Retry attempt {attempt + 1} after error: {e}")
def say_hello_with_metadata(self, name, metadata=None):
"""带元数据的RPC调用"""
if metadata is None:
metadata = [('client-id', 'python-client-1')]
request = helloworld_pb2.HelloRequest(name=name)
try:
response, call = self.stub.SayHello.with_call(
request,
metadata=metadata,
timeout=self.timeout
)
# 获取响应元数据
trailing_metadata = call.trailing_metadata()
print(f"Response metadata: {trailing_metadata}")
return response.message
except grpc.RpcError as e:
print(f"RPC failed: {e.code()}")
return None
def say_hello_async(self, name, callback: Optional[Callable] = None):
"""异步调用"""
request = helloworld_pb2.HelloRequest(name=name)
# 创建异步调用
future = self.stub.SayHello.future(request, timeout=self.timeout)
if callback:
# 添加回调
future.add_done_callback(
lambda f: callback(f.result().message) if not f.exception()
else print(f"Error: {f.exception()}")
)
return future
def get_stats(self):
"""获取通道统计信息"""
stats = []
# 连接状态
channel_connectivity = self.channel.get_state(try_to_connect=True)
stats.append(f"Channel state: {channel_connectivity}")
return stats
def close(self):
"""优雅关闭"""
self.channel.close()
# 高级使用示例
def advanced_usage():
client = AdvancedGRPCClient(
host='server.example.com',
port=443,
use_ssl=True,
ssl_cert='server.crt',
max_retries=3,
timeout=5
)
try:
# 带元数据的调用
result = client.say_hello_with_metadata(
'Alice',
metadata=[('auth-token', 'my-token')]
)
print(f"Result: {result}")
# 异步调用
future = client.say_hello_async('Bob', callback=lambda msg: print(f"Async result: {msg}"))
result = future.result() # 等待完成
print(f"Async result: {result}")
# 获取统计信息
stats = client.get_stats()
for stat in stats:
print(stat)
finally:
client.close()
if __name__ == '__main__':
advanced_usage()
客户端配置选项
import grpc
from grpc import aio
import asyncio
class ConfigurableClient:
@staticmethod
def create_secure_channel(host, port, cert_file, key_file=None):
"""创建安全的gRPC通道"""
with open(cert_file, 'rb') as f:
trusted_certs = f.read()
if key_file:
# 双向TLS
with open(key_file, 'rb') as f:
private_key = f.read()
credentials = grpc.ssl_channel_credentials(
root_certificates=trusted_certs,
private_key=private_key,
certificate_chain=trusted_certs
)
else:
# 服务端TLS
credentials = grpc.ssl_channel_credentials(trusted_certs)
return grpc.secure_channel(f'{host}:{port}', credentials)
@staticmethod
def create_channel_with_options(host, port):
"""带选项的通道创建"""
options = [
('grpc.max_send_message_length', 50 * 1024 * 1024), # 50MB
('grpc.max_receive_message_length', 50 * 1024 * 1024), # 50MB
('grpc.keepalive_time_ms', 10000),
('grpc.keepalive_timeout_ms', 5000),
('grpc.enable_retries', 1),
('grpc.service_config', json.dumps({
"methodConfig": [{
"name": [{"service": "helloworld.Greeter"}],
"retryPolicy": {
"maxAttempts": 4,
"initialBackoff": "0.1s",
"maxBackoff": "1s",
"backoffMultiplier": 2,
"retryableStatusCodes": ["UNAVAILABLE"]
}
}]
}))
]
return grpc.insecure_channel(f'{host}:{port}', options=options)
# 异步客户端
class AsyncGRPCClient:
def __init__(self, host='localhost', port=50051):
self.host = host
self.port = port
self.channel = None
self.stub = None
async def connect(self):
"""异步连接"""
self.channel = aio.insecure_channel(f'{self.host}:{self.port}')
self.stub = helloworld_pb2_grpc.GreeterStub(self.channel)
async def say_hello(self, name):
"""异步调用"""
request = helloworld_pb2.HelloRequest(name=name)
response = await self.stub.SayHello(request)
return response.message
async def close(self):
"""关闭连接"""
await self.channel.close()
# 异步使用示例
async def async_main():
client = AsyncGRPCClient()
await client.connect()
try:
result = await client.say_hello('AsyncUser')
print(f"Async result: {result}")
finally:
await client.close()
if __name__ == '__main__':
asyncio.run(async_main())
错误处理与健康检查
import grpc
from grpc_health.v1 import health_pb2, health_pb2_grpc
class RobustClient:
def __init__(self, host='localhost', port=50051):
self.channel = grpc.insecure_channel(f'{host}:{port}')
self.stub = helloworld_pb2_grpc.GreeterStub(self.channel)
self.health_stub = health_pb2_grpc.HealthStub(self.channel)
def check_health(self, service='helloworld.Greeter'):
"""健康检查"""
try:
request = health_pb2.HealthCheckRequest(service=service)
response = self.health_stub.Check(request)
return response.status == health_pb2.HealthCheckResponse.SERVING
except Exception as e:
print(f"Health check failed: {e}")
return False
def robust_say_hello(self, name, max_retries=3):
"""健壮的RPC调用"""
last_error = None
for attempt in range(max_retries):
try:
if self.check_health():
response = self.stub.SayHello(
helloworld_pb2.HelloRequest(name=name),
timeout=5
)
return response.message
else:
print("Service not healthy, waiting...")
time.sleep(1)
except grpc.RpcError as e:
last_error = e
if e.code() == grpc.StatusCode.UNAVAILABLE:
print(f"Service unavailable, retrying... (attempt {attempt + 1})")
time.sleep(2 ** attempt) # 指数退避
else:
raise
raise last_error or Exception("Max retries exceeded")
这个完整的实现包含了:
- 四种RPC类型(一元、服务端流、客户端流、双向流)
- SSL/TLS安全连接
- 元数据和认证
- 异步和同步调用
- 重试和错误处理
- 健康检查
- 通道配置优化
根据你的具体需求选择合适的功能模块。