gRPC基于HTTP/2和Protocol Buffers实现高性能RPC通信,在微服务架构中相比REST+JSON方案具有更低的序列化开销和更强的类型约束。后端开发场景中,内部服务间通信的效率和可靠性直接影响系统整体性能。本文演示从Protobuf定义、服务实现到客户端调用的完整gRPC微服务开发流程。
Protobuf接口定义与代码生成
Protocol Buffers作为gRPC的接口定义语言(IDL),通过proto文件描述服务接口和消息结构。proto文件同时是服务端和客户端的契约,代码生成工具据此生成多语言绑定。
// proto/user_service.proto
syntax = "proto3";
package userservice.v1;
option go_package = "github.com/example/proto/userservice/v1";
option java_package = "com.example.proto.userservice.v1";
// 用户服务定义
service UserService {
// 创建用户
rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
// 获取用户信息
rpc GetUser(GetUserRequest) returns (GetUserResponse);
// 批量获取用户(服务端流式)
rpc ListUsers(ListUsersRequest) returns (stream User);
// 双向流式聊天
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
message User {
string id = 1;
string name = 2;
string email = 3;
int32 age = 4;
UserStatus status = 5;
int64 created_at = 6;
}
enum UserStatus {
USER_STATUS_UNSPECIFIED = 0;
USER_STATUS_ACTIVE = 1;
USER_STATUS_INACTIVE = 2;
USER_STATUS_BANNED = 3;
}
message CreateUserRequest {
string name = 1;
string email = 2;
int32 age = 3;
}
message CreateUserResponse {
User user = 1;
}
message GetUserRequest {
string id = 1;
}
message GetUserResponse {
User user = 1;
}
message ListUsersRequest {
int32 page_size = 1;
string page_token = 2;
UserStatus status_filter = 3;
}
message ChatMessage {
string user_id = 1;
string content = 2;
int64 timestamp = 3;
}
# 安装工具链
pip install grpcio grpcio-tools protobuf
# 生成Python代码
python -m grpc_tools.protoc \
--proto_path=proto \
--python_out=src/generated \
--grpc_python_out=src/generated \
proto/user_service.proto
# 生成的文件结构
# src/generated/
# ├── user_service_pb2.py # 消息类定义
# └── user_service_pb2_grpc.py # gRPC服务stub和servicer
gRPC服务端实现与拦截器配置
服务端实现UserServiceServicer接口,处理客户端请求。通过拦截器实现认证、日志、指标采集等横切关注点。
# src/server.py
import grpc
import time
import logging
from concurrent import futures
from generated import user_service_pb2 as pb
from generated import user_service_pb2_grpc as pb_grpc
logger = logging.getLogger(__name__)
class UserServiceServicer(pb_grpc.UserServiceServicer):
def __init__(self, db_session):
self.db = db_session
def CreateUser(self, request, context):
# 参数校验
if not request.name or not request.email:
context.abort(grpc.StatusCode.INVALID_ARGUMENT,
"name和email不能为空")
# 检查邮箱是否已存在
existing = self.db.query(User).filter_by(email=request.email).first()
if existing:
context.abort(grpc.StatusCode.ALREADY_EXISTS,
"邮箱已被注册")
user = User(
name=request.name,
email=request.email,
age=request.age,
status=pb.UserStatus.USER_STATUS_ACTIVE
)
self.db.add(user)
self.db.commit()
return pb.CreateUserResponse(
user=pb.User(
id=str(user.id),
name=user.name,
email=user.email,
age=user.age,
status=user.status,
created_at=int(time.time())
)
)
def GetUser(self, request, context):
user = self.db.query(User).filter_by(id=request.id).first()
if not user:
context.abort(grpc.StatusCode.NOT_FOUND, "用户不存在")
return pb.GetUserResponse(
user=pb.User(
id=str(user.id),
name=user.name,
email=user.email,
status=user.status
)
)
# 服务端流式 - 批量返回用户
def ListUsers(self, request, context):
query = self.db.query(User)
if request.status_filter:
query = query.filter(User.status == request.status_filter)
offset = 0
while True:
batch = query.offset(offset).limit(request.page_size).all()
if not batch:
break
for user in batch:
yield pb.User(
id=str(user.id),
name=user.name,
email=user.email
)
offset += request.page_size
# 双向流式
def Chat(self, request_iterator, context):
for message in request_iterator:
yield pb.ChatMessage(
user_id="system",
content=f"收到来自{message.user_id}的消息: {message.content}",
timestamp=int(time.time())
)
# 认证拦截器
class AuthInterceptor(grpc.ServerInterceptor):
def __init__(self, valid_tokens):
self.valid_tokens = valid_tokens
def intercept_service(self, continuation, handler_call_details):
metadata = dict(handler_call_details.invocation_metadata)
token = metadata.get('authorization')
if token not in self.valid_tokens:
return grpc.unary_unary_rpc_method_handler(
lambda request, context: context.abort(
grpc.StatusCode.UNAUTHENTICATED, "无效的认证token"
)
)
return continuation(handler_call_details)
# 日志拦截器
class LoggingInterceptor(grpc.ServerInterceptor):
def intercept_service(self, continuation, handler_call_details):
start = time.time()
method = handler_call_details.method
try:
result = continuation(handler_call_details)
elapsed = time.time() - start
logger.info(f"{method} - {elapsed:.3f}s")
return result
except Exception as e:
elapsed = time.time() - start
logger.error(f"{method} - {elapsed:.3f}s - ERROR: {e}")
raise
def serve():
server = grpc.server(
futures.ThreadPoolExecutor(max_workers=50),
interceptors=[
AuthInterceptor(valid_tokens={"secret-token-123"}),
LoggingInterceptor()
],
options=[
('grpc.max_send_message_length', 100 * 1024 * 1024),
('grpc.max_receive_message_length', 100 * 1024 * 1024),
('grpc.keepalive_time', 30000),
('grpc.keepalive_timeout', 10000),
]
)
pb_grpc.add_UserServiceServicer_to_server(
UserServiceServicer(get_db_session()),
server
)
server.add_insecure_port('[::]:50051')
server.start()
logger.info("gRPC server started on port 50051")
server.wait_for_termination()
if __name__ == '__main__':
serve()
gRPC客户端调用与连接池管理
客户端通过channel连接服务端,支持Unary、服务端流式、客户端流式和双向流式四种调用模式。生产环境需管理连接池和重试策略。
# src/client.py
import grpc
import json
import time
from generated import user_service_pb2 as pb
from generated import user_service_pb2_grpc as pb_grpc
# 重试拦截器
class RetryInterceptor(grpc.UnaryUnaryClientInterceptor):
def __init__(self, max_retries=3):
self.max_retries = max_retries
def intercept_unary_unary(self, continuation, client_call_details, request):
for attempt in range(self.max_retries):
try:
response = continuation(client_call_details, request)
return response
except grpc.RpcError as e:
if e.code() in (grpc.StatusCode.UNAVAILABLE,
grpc.StatusCode.DEADLINE_EXCEEDED):
if attempt < self.max_retries - 1:
time.sleep(0.5 * (attempt + 1))
continue
raise
return None
# 连接配置
channel = grpc.insecure_channel(
'localhost:50051',
options=[
('grpc.max_send_message_length', 100 * 1024 * 1024),
('grpc.max_receive_message_length', 100 * 1024 * 1024),
('grpc.keepalive_time', 30000),
('grpc.keepalive_timeout', 10000),
('grpc.enable_retries', 1),
('grpc.service_config', json.dumps({
"methodConfig": [{
"name": [{"service": "userservice.v1.UserService"}],
"retryPolicy": {
"maxAttempts": 3,
"initialBackoff": "0.1s",
"maxBackoff": "1s",
"multiplier": 2,
"retryableStatusCodes": [
"UNAVAILABLE", "DEADLINE_EXCEEDED"
]
}
}]
}))
],
interceptors=[RetryInterceptor(max_retries=3)]
)
stub = pb_grpc.UserServiceStub(channel)
# Unary调用
response = stub.CreateUser(
pb.CreateUserRequest(name="张三", email="zhangsan@example.com", age=28),
metadata=[('authorization', 'secret-token-123')],
timeout=5.0
)
print(f"创建用户成功: {response.user.id}")
# 服务端流式调用
for user in stub.ListUsers(
pb.ListUsersRequest(page_size=10, status_filter=pb.UserStatus.USER_STATUS_ACTIVE),
metadata=[('authorization', 'secret-token-123')]
):
print(f"用户: {user.name} - {user.email}")
# 双向流式调用
def generate_messages():
for i in range(5):
yield pb.ChatMessage(
user_id="client-001",
content=f"消息{i}",
timestamp=int(time.time())
)
for response in stub.Chat(generate_messages(),
metadata=[('authorization', 'secret-token-123')]):
print(f"收到回复: {response.content}")
健康检查与负载均衡集成
gRPC内置健康检查协议,与Kubernetes和Envoy等服务网格集成时,通过健康检查协议实现服务发现和流量切换。
# src/health.py - 健康检查服务
import grpc
from grpc_health.v1 import health, health_pb2, health_pb2_grpc
def setup_health_check(server, service_name="userservice.v1.UserService"):
health_servicer = health.HealthCheckServicer()
health_servicer.set(service_name, health_pb2.HealthCheckResponse.SERVING)
health_pb2_grpc.add_HealthServicer_to_server(health_servicer, server)
return health_servicer
# Kubernetes gRPC健康检查探针配置
# livenessProbe:
# grpc:
# port: 50051
# initialDelaySeconds: 10
# periodSeconds: 10
# readinessProbe:
# grpc:
# port: 50051
# initialDelaySeconds: 5
# periodSeconds: 5
生产环境建议使用TLS加密通信(grpc.secure_channel配合SSL证书),内部微服务通信可使用mTLS双向认证。gRPC与REST可通过grpc-gateway自动生成HTTP代理,对外暴露REST接口同时内部使用gRPC。连接管理方面,长连接场景需配置keepalive参数防止NAT超时断连,短连接场景需注意channel复用以避免频繁建连开销。对于高并发调用,建议使用grpc.aio异步客户端(基于asyncio),提升I/O利用率。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/pythongrpc-wei-fu-wu-tong-xin-yu-protobuf-jie-kou-ding-yi/