Erlang/Elixir gRPC 与 Protobuf:从编码原理到 grpcbox 实战

系统讲解 BEAM 生态的 gRPC 技术栈:Protobuf 二进制编码原理与 .proto 语法、gpb 代码生成与 Erlang 集成、grpcbox 服务端与客户端实现、四种流式 RPC 语义、拦截器与元数据传递,以及超时、状态码与生产实践。

在微服务架构里,服务间的契约描述往往比服务本身的实现更容易腐化。REST + JSON 的组合虽然上手快,却缺少强制的接口定义:字段拼错、类型不符、必填项遗漏这些问题只会在运行时暴露,而且 JSON 文本编码的体积与解析开销在高频调用下相当可观。gRPC 用 Protocol Buffers 作为接口定义语言(IDL)与序列化格式,把「接口契约」变成编译期产物——服务端与客户端从同一份 .proto 生成代码,字段类型不匹配根本编译不过。

对 BEAM 生态而言,gRPC 的吸引力还多一层:Erlang 与 Elixir 的进程模型天然适合处理长连接与流式 RPC,一个 gRPC 流就是一组轻量进程之间的消息传递,无需线程池、无需回调地狱。本文从 Protobuf 的二进制编码讲起,逐层展开 gpb 代码生成、grpcbox 服务端与客户端实现、四种流式语义与拦截器,最后给出生产环境的状态码与超时治理实践。

一、为什么需要 gRPC 与 Protobuf

1.1 与 REST/JSON 的对比

维度REST + JSONgRPC + Protobuf
接口定义OpenAPI(可选,易与实现脱节).proto(强制,编译期校验)
编码文本 JSON,字段名重复传输二进制,仅传字段号
体积基准通常为 JSON 的 20%~40%
解析速度需词法+语法分析按偏移直接读取
传输HTTP/1.1(多为短连接)HTTP/2(多路复用长连接)
流式SSE/WebSocket 需另建协议原生四种流式语义
浏览器支持原生需 grpc-web 代理

1.2 什么时候不该用 gRPC

  • 面向浏览器的公开 API:浏览器无法直接调用 gRPC,需要 grpc-web 代理层,增加一跳;
  • 对外合作伙伴接口:对方未必有 Protobuf 工具链,JSON 的通用性更稳妥;
  • 低频、调试优先的内部接口:curl 能直接调的 REST 在排障时优势明显。

gRPC 的最佳战场是内部服务间的高频调用:东西向流量、强类型契约、需要流式或双向通信的场景。

1.3 BEAM 上的实现选型

项目语言特点
grpcboxErlang官方 grpc 库之上的 OTP 封装,基于 Cowboy,适合 Erlang 项目
gpbErlangProtobuf 代码生成器,输出纯 Erlang 模块
elixir-grpcElixirGRPC.Server/GRPC.Stub 宏,Phoenix 生态集成好
protobufElixirElixir 版 Protobuf 编解码库

Erlang 项目选 grpcbox + gpb,Elixir 项目选 elixir-grpc + protobuf,两者可以互通——线格式由 Protobuf 与 HTTP/2 标准定义,与实现语言无关。

二、Protobuf 编码原理与 .proto 语法

2.1 线格式:字段号 + 线类型

Protobuf 不传字段名,只传「字段号 + 线类型」,这正是它比 JSON 紧凑的根本原因。每个字段的编码以 tag 开头,tag = (field_number << 3) | wire_type:

线类型值承载适用类型
Varint0变长整数int32/int64/uint/bool/enum
64-bit1固定 8 字节fixed64/sfixed64/double
Length-delimited2长度前缀 + 数据string/bytes/嵌套消息/repeated
32-bit5固定 4 字节fixed32/sfixed32/float

举例:字段号 1、类型 string 的 name 字段,tag 为 (1 << 3) | 2 = 0x0A。字段号 1~15 只需一个字节的 tag,因此高频字段应分配小字段号。

2.2 Varint 与 ZigZag

Varint 用每个字节的最高位表示「是否还有后续字节」,小数值只占一字节:300 编码为 0xAC 0x02 两个字节。负数若直接用 Varint 会占满 10 字节(补码全 1),所以 sint32/sint64 采用 ZigZag 映射,把负数挤到小数值区间:

ZigZag(n) = (n << 1) ^ (n >> 31)
  0 → 0     -1 → 1     1 → 2     -2 → 3

结论:可能为负的整数字段,一律用 sint32/sint64,否则体积会翻几倍。

2.3 一个完整的 .proto 定义

syntax = "proto3";

package user.v1;

import "google/protobuf/timestamp.proto";

service UserService {
  rpc GetUser(GetUserRequest) returns (User);
  rpc ListUsers(ListUsersRequest) returns (stream User);
  rpc UploadUsers(stream User) returns (UploadSummary);
  rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message GetUserRequest { int64 user_id = 1; }

message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  repeated string roles = 4;
  google.protobuf.Timestamp created_at = 5;
  map<string, string> labels = 6;
  optional string nickname = 7;
}

message UploadSummary { int32 accepted = 1; int32 rejected = 2; }

2.4 proto3 的关键语义

特性说明
默认值不传输字段等于零值时不上线,接收方读到默认值
optional显式存在性检测,生成 has_xxx 语义
repeated数值类型可加 [packed=true] 压缩编码
map语法糖,线格式等价于 repeated MapEntry
字段号不可变一旦上线,字段号与类型都不能改,只能新增
reserved删除字段时保留其编号,防止被误复用

字段号复用是 Protobuf 最危险的陷阱:删掉字段 4 后又用它定义新类型,老客户端会把新字段按旧类型解析,得到静默错误的数据。

三、gpb 代码生成与 Erlang 集成

3.1 rebar3 项目配置

%% rebar.config
{plugins, [rebar3_gpb_plugin]}.

{gpb_opts, [
    {i, "priv/proto"},                  %% 导入路径
    {module_name_suffix, "_pb"},        %% 生成 user_pb 而非 user
    {strings, as_binary},               %% string 映射为 binary(推荐)
    {maps, true},                       %% 消息映射为 maps 而非 record
    {type_specs, true}                  %% 生成 -spec 供 Dialyzer 使用
]}.

strings as_binary 与 maps true 是现代 Erlang 项目的标配:前者避免 list 带来的 GC 压力,后者让模式匹配写起来更直观。

3.2 生成产物

rebar3 compile 会在 _build 下生成 user_pb.erl,核心导出函数如下:

函数作用
encode_msg(Msg, Module)把 map/record 编码为 iolist
decode_msg(Bin, Module)从二进制解码为 map
verify_msg(Msg, Module)校验消息结构,返回 ok 或错误
get_msg_defs()返回消息定义(供反射/校验)
Msg = #{id => 42, name => <<"Alice">>, email => <<"a@b.c">>, roles => [<<"admin">>]},
Bin = user_pb:encode_msg(Msg, 'user.v1.User'),
Decoded = user_pb:decode_msg(iolist_to_binary(Bin), 'user.v1.User').

3.3 与 Ecto/数据库结构的映射

Protobuf 的 map 与 Elixir 的 struct 之间需要一个显式的转换层。不要试图用宏自动映射——字段名与类型差异会在重构时变成灾难。手写 to_proto/1 与 from_proto/1 虽然啰嗦,但可读且可测:

defmodule MyApp.Proto.UserMapper do
  alias User.V1.User

  def to_proto(%MyApp.Accounts.User{} = u) do
    %User{id: u.id, name: u.name, email: u.email, roles: u.roles,
          created_at: to_ts(u.inserted_at), labels: u.labels || %{}}
  end

  defp to_ts(%DateTime{} = dt),
    do: Google.Protobuf.Timestamp.new(seconds: DateTime.to_unix(dt))
end

四、grpcbox 服务端实现

4.1 定义服务行为

-module(user_service).
-behaviour(grpcbox_service).

-export([get_user/2, list_users/2]).

%% 一元 RPC:请求 -> 响应
get_user(GetUserReq, Stream) ->
    #{user_id := Id} = GetUserReq,
    case user_db:find(Id) of
        {ok, User} -> {ok, to_pb(User), Stream};
        error      -> {grpcbox:status(not_found, <<"user not found">>), Stream}
    end.

%% 服务端流式:一次请求,多次响应
list_users(#{page_size := Size}, Stream) ->
    Users = user_db:list(Size),
    lists:foldl(fun(U, S) -> grpcbox:send_response(to_pb(U), S) end, Stream, Users),
    {ok, Stream}.

4.2 启动配置

%% sys.config
{grpcbox, [
    {servers, [
        #{grpc_opts => #{
              service_protos => [user_pb],
              services => #{'user.v1.UserService' => user_service}
          },
          listen_opts => #{port => 8080, ip => {0, 0, 0, 0}},
          transport_opts => #{max_connections => 10000, idle_timeout => 30000}}
    ]}
]}.

grpcbox 默认使用 Cowboy 作为 HTTP/2 传输层,因此它也继承了 Cowboy 的连接管理能力(见 https://plumephp.com/erlang-cowboy-http-api/)。

4.3 客户端流式

upload_users(Stream, _InitState) ->
    upload_loop(Stream, 0, 0).

upload_loop(Stream, Accepted, Rejected) ->
    case grpcbox:recv_request(Stream) of
        {ok, User} ->
            case validate(User) of
                ok      -> upload_loop(Stream, Accepted + 1, Rejected);
                invalid -> upload_loop(Stream, Accepted, Rejected + 1)
            end;
        eof ->
            {ok, #{accepted => Accepted, rejected => Rejected}, Stream}
    end.

双向流的回调更复杂,需要同时处理「收到消息」与「可以发送」两类事件,通常用 handle_info/2 配合状态机实现——本质上是把 gRPC 流当作一个 gen_statem 来驱动。

4.4 四种 RPC 语义对照

语义服务端回调适用场景
一元返回 {ok, Resp, Stream}查询、下单等常规调用
服务端流循环 send_response 后返回 {ok, Stream}大结果集分页推送、日志跟随
客户端流循环 recv_request 直到 eof批量上传、指标上报
双向流handle_info 驱动状态机实时协作、聊天、长连接协商

五、grpcbox 客户端、流式 RPC 与拦截器

5.1 建立连接与一元调用

{ok, Channel} = grpcbox_channel:start_link(user_channel, #{
    endpoints => [{http, "user-svc.internal", 8080, []}],
    pool_size => 4
}),

{ok, Reply} = grpcbox_client:unary(
    Channel, <<"/user.v1.UserService/GetUser">>, #{user_id => 42},
    #{timeout => 5000}),
Name = maps:get(name, Reply).

5.2 服务端流式调用

{ok, Stream} = grpcbox_client:stream(
    Channel, <<"/user.v1.UserService/ListUsers">>, #{page_size => 100},
    #{timeout => 30000}),

collect(Stream).
collect(Stream) ->
    case grpcbox_client:recv(Stream) of
        {ok, Msg} -> [Msg | collect(Stream)];
        eof       -> []
    end.

5.3 元数据传递

元数据是 gRPC 的「HTTP 头」,用于传递认证信息、链路追踪 ID 等非业务数据:

%% 客户端发送元数据
MD = #{<<"authorization">> => <<"Bearer ", Token/binary>>,
       <<"x-request-id">> => RequestId},
grpcbox_client:unary(Channel, Path, Req, #{metadata => MD, timeout => 5000}).

%% 服务端读取元数据
get_user(Req, Stream) ->
    case maps:get(<<"authorization">>, grpcbox:metadata(Stream), undefined) of
        undefined -> {grpcbox:status(unauthenticated, <<"missing token">>), Stream};
        Token     -> handle_with_token(Token, Req, Stream)
    end.

5.4 拦截器

拦截器是横切关注点(认证、日志、追踪、限流)的统一入口:

-module(auth_interceptor).
-behaviour(grpcbox_interceptor).

-export([unary/5, stream/4]).

unary(_Path, Req, Stream, MD, _Opts) ->
    case verify(MD) of
        {ok, Claims} -> grpcbox:continue(Req, Stream, MD#{claims => Claims});
        {error, R}   -> grpcbox:abort(Stream, #{code => unauthenticated, message => R})
    end.

stream(_Path, Stream, MD, _Opts) ->
    case verify(MD) of
        {ok, _}    -> grpcbox:continue(Stream, MD);
        {error, R} -> grpcbox:abort(Stream, #{code => unauthenticated, message => R})
    end.

注册顺序与执行顺序一致:interceptors => [auth_interceptor, trace_interceptor, metrics_interceptor]。认证应排在追踪与限流之前——被拒绝的请求不该消耗限流配额。

六、超时、错误码与生产实践

6.1 标准状态码

码名称典型场景
0OK成功
3INVALID_ARGUMENT参数校验失败(不重试)
4DEADLINE_EXCEEDED超时(可重试,需幂等)
5NOT_FOUND资源不存在
7PERMISSION_DENIED已认证但无权限
8RESOURCE_EXHAUSTED限流/配额耗尽(应退避)
13INTERNAL服务内部错误(需告警)
14UNAVAILABLE服务不可达(重试 + 熔断)
16UNAUTHENTICATED未认证或凭证失效

重试策略与状态码强相关:只有 UNAVAILABLE、DEADLINE_EXCEEDED、RESOURCE_EXHAUSTED 值得重试,且必须配合指数退避与幂等保证;INVALID_ARGUMENT 与 NOT_FOUND 重试多少次都是徒劳。

6.2 超时传播

超时应当从入口一路传到最底层,且每跳都预留余量:

客户端 deadline 1000ms
   └─▶ 网关 剩余 900ms ──▶ 服务 A 剩余 700ms ──▶ 服务 B 剩余 400ms
        (服务 B 留给自己的处理时间只有 300ms)

gRPC 的 grpc-timeout 头会自动随请求传播,服务端可以用 grpcbox:deadline(Stream) 读取剩余时间,据此提前放弃无望的调用,避免做「已经超时的无用功」。

6.3 连接池与负载均衡

{ok, Channel} = grpcbox_channel:start_link(user_channel, #{
    endpoints => [{http, "user-svc-1.internal", 8080, []},
                  {http, "user-svc-2.internal", 8080, []}],
    pool_size => 8,
    load_balancing => round_robin
}),

在 Kubernetes 环境中,gRPC 的负载均衡有一个经典陷阱:服务发现返回的是 ClusterIP,HTTP/2 长连接建立后所有请求都打到同一个 Pod。解法是使用 headless Service 拿到所有 Pod IP,或在客户端接入 xDS/服务网格。

6.4 与可观测性结合

gRPC 的元数据是分布式追踪的天然载体。trace_interceptor 从入站元数据中提取 traceparent(W3C 标准),生成 span 后再注入到出站调用:

stream(Path, Stream, Metadata, _Opts) ->
    ParentCtx = otel_propagator:text_map_extract(Metadata),
    Ctx = otel_tracer:start_span(Path, #{parent => ParentCtx}),
    grpcbox:continue(Stream, Metadata#{otel_ctx => Ctx}).

配套指标至少应包含每个方法的调用次数、P50/P99 延迟、状态码分布,与 https://plumephp.com/erlang-logging-telemetry-observability/ 中描述的 Telemetry 体系共用同一套采集管道。

6.5 版本兼容演进

Protobuf 的向后兼容规则必须严格遵守,否则滚动升级期间新旧版本会互相解析失败:

  • 只增不改:新增字段永远用新的字段号,绝不修改已有字段的类型或编号;
  • 删除用 reserved:reserved 4, 6; reserved "old_field";,防止编号被复用;
  • enum 加 UNSPECIFIED = 0:proto3 要求 enum 第一个值为 0,客户端遇到未知枚举值时应优雅降级而非报错;
  • 服务方法只增不删:删除方法会让老客户端收到 UNIMPLEMENTED,如需下线应先用拦截器统计调用量,确认为零后再移除。

七、最佳实践与总结

  • .proto 独立成仓库或用 git submodule:契约是跨服务的公共资产,放在某个服务的代码库里会导致其他服务被迫依赖该服务的发布节奏;
  • 字段号规划留空:核心消息的字段号从 1 开始连续分配,同时用 reserved 预留区间给未来的高频字段,避免后期插入破坏小 tag 的紧凑优势;
  • string 用 binary,负数用 sint:Erlang 侧 {strings, as_binary},proto 侧可负整数字段一律 sint32/sint64;
  • 拦截器顺序有讲究:认证 → 追踪 → 限流 → 日志,被拒绝的请求不应消耗配额与产生噪音;
  • 状态码要语义正确:把「参数错误」也返回 INTERNAL,会让客户端的重试策略彻底失效;
  • 超时逐跳收敛:每一跳都要留出自己的处理余量,否则会出现「上游超时了,下游还在干活」的资源浪费;
  • 连接池与 LB 单独验证:在 K8s 上部署后务必确认流量是否真的分散到所有 Pod,这是 gRPC 上生产最容易踩的坑;
  • 契约变更走兼容性检查:把 buf breaking 之类的工具接入 CI,自动拦截破坏性变更。

gRPC 与 Protobuf 的价值不在「性能更好」,而在于把接口契约从文档变成了编译产物。当服务数量超过十个,跨团队沟通的成本会迅速超过编码本身,此时一份强制的、可演进的、可自动校验的 IDL 就是最划算的投资。而 BEAM 的进程模型让 gRPC 的流式语义实现起来异常自然——一条流就是一组进程,背压来自消息队列长度,容错来自监督树。把这两者结合,你就得到了一个既强类型又高容错的内部通信层。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「erlang」更多文章

  1. Erlang/Elixir 容器化与集群部署:Release、Docker 与 libcluster
  2. Elixir 认证授权实战:JWT、Guardian 与 Phoenix.Token
  3. Elixir HTTP 客户端与连接池:Mint、Finch 与 Req 实战