Python脚本gRPC客户端如何实现

wen 实用脚本 5

本文目录导读:

Python脚本gRPC客户端如何实现

  1. 环境准备
  2. 定义Proto文件
  3. 生成Python代码
  4. 基础客户端实现
  5. 高级客户端实现
  6. 客户端配置选项
  7. 错误处理与健康检查

我来详细介绍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安全连接
  • 元数据和认证
  • 异步和同步调用
  • 重试和错误处理
  • 健康检查
  • 通道配置优化

根据你的具体需求选择合适的功能模块。

抱歉,评论功能暂时关闭!