Cowboy 与 HTTP JSON API:从路由到实时通信

系统讲解 Cowboy 的架构与请求生命周期、路由与 handler、中间件与 JSON 编解码、WebSocket 与 SSE 实时通信、限流与安全加固以及生产部署。

Cowboy 是 Erlang 生态中最主流的 HTTP 服务器,也是 Elixir Phoenix 底层的 HTTP 层。它完全基于 OTP 构建:每个连接是一个进程,每个请求在连接进程内被 handler 处理,因此天然继承了 BEAM 的隔离性与容错能力。Cowboy 同时支持 HTTP/1.1、HTTP/2、WebSocket 与 SSE,是构建 JSON API 与实时服务的理想底座。本文将从请求生命周期讲起,覆盖路由、handler、中间件、JSON 编解码、实时通信、限流安全与部署,给出一套可直接落地的 API 服务骨架。

一、Cowboy 架构与请求生命周期

1.1 连接即进程

Cowboy 的并发模型非常直接:每个 TCP 连接由一个 Erlang 进程处理。

Listener (ranch)
└── Acceptor Pool
    └── Connection Process (每连接一个)
        ├── 解析 HTTP 请求
        ├── 依次执行 middleware
        ├── 调用 handler 的 init/2
        └── 发送响应后进入 keep-alive 或关闭

这带来两个直接好处:一个连接的崩溃不影响其他连接;可以轻松支撑数十万并发连接(每个连接进程初始内存仅几 KB)。

1.2 请求生命周期

  1. Ranch 接受 TCP 连接,创建连接进程;
  2. 连接进程解析请求行与头部,构造 Req 对象;
  3. 依次执行配置的 middleware(如 cowboy_router、cowboy_handler);
  4. 路由匹配到 handler 模块,调用 Handler:init(Req, Opts);
  5. handler 返回 {ok, Req, State},Cowboy 发送响应。

1.3 启动一个 Cowboy 服务

-module(api_app).
-behaviour(application).
-export([start/2, stop/1]).

start(_Type, _Args) ->
    Dispatch = cowboy_router:compile([
        {'_', [
            {"/api/users", user_handler, #{}},
            {"/api/users/:id", user_handler, #{}},
            {"/api/health", health_handler, #{}},
            {"/ws", ws_handler, #{}},
            {"/[...]", not_found_handler, #{}}
        ]}
    ]),
    {ok, _} = cowboy:start_clear(http_listener,
        [{port, 8080}, {num_acceptors, 10}, {max_connections, 100_000}],
        #{env => #{dispatch => Dispatch},
          middlewares => [cowboy_router, cowboy_handler],
          idle_timeout => 60_000,
          request_timeout => 5_000}),
    api_sup:start_link().

stop(_State) -> ok.

二、路由与 Handler 实现

2.1 路由规则

Dispatch = cowboy_router:compile([
    {'_', [                                          % '_' 匹配任意 Host
        {"/api/users",        user_list_h,   #{}},
        {"/api/users/:id",    user_item_h,   #{}},   % :id 绑定变量
        {"/api/users/:id/posts/:post_id", user_post_h, #{}},
        {"/static/[...]",     cowboy_static, {priv_dir, my_app, "static"}},
        {"/[...]",            not_found_h,   #{}}
    ]}
]).
模式含义
/api/users精确匹配
/api/users/:id:id 绑定到 bindings
/static/[...]前缀匹配,剩余部分进 path_info
/[...]兜底捕获所有路径

2.2 Handler 的 init/2

所有 handler 实现 init/2 回调,按 HTTP method 分派:

-module(user_item_h).
-export([init/2]).

init(Req0, State) ->
    Method = cowboy_req:method(Req0),
    Id = cowboy_req:binding(id, Req0),
    Req = handle(Method, Id, Req0),
    {ok, Req, State}.

handle(<<"GET">>, Id, Req) ->
    case db:find_user(Id) of
        {ok, User} ->
            reply_json(200, user_to_json(User), Req);
        {error, not_found} ->
            reply_error(404, <<"user not found">>, Req)
    end;
handle(<<"DELETE">>, Id, Req) ->
    ok = db:delete_user(Id),
    cowboy_req:reply(204, #{}, <<>>, Req);
handle(_, _, Req) ->
    reply_error(405, <<"method not allowed">>, Req).

reply_json(Code, Data, Req) ->
    cowboy_req:reply(Code,
        #{<<"content-type">> => <<"application/json; charset=utf-8">>,
          <<"cache-control">> => <<"no-store">>},
        jsx:encode(Data), Req).

reply_error(Code, Msg, Req) ->
    reply_json(Code, #{error => Msg, code => Code}, Req).

2.3 读取请求体与查询参数

%% 读取查询参数与单个头部
Qs = cowboy_req:parse_qs(Req),
Page = proplists:get_value(<<"page">>, Qs, <<"1">>),
Auth = cowboy_req:header(<<"authorization">>, Req),

%% 读取请求体(必须设置长度上限,防 OOM)
read_body_limited(Req, Max) ->
    case cowboy_req:read_body(Req, #{length => Max, period => 5000}) of
        {ok, Body, Req2} -> {ok, Body, Req2};
        {more, _Partial, _Req2} -> {error, payload_too_large}
    end.

2.4 流式响应

%% 大文件或分块生成:先发头部,再逐块推送
Req2 = cowboy_req:stream_reply(200, #{<<"content-type">> => <<"text/plain">>}, Req),
cowboy_req:stream_body(<<"chunk 1\n">>, nofin, Req2),
cowboy_req:stream_body(<<"chunk 2\n">>, fin, Req2).

三、中间件与 JSON 编解码

3.1 中间件机制

中间件是 Cowboy 处理链的扩展点,可插入日志、认证、CORS、请求 ID 等横切逻辑:

-module(auth_middleware).
-behaviour(cowboy_middleware).
-export([execute/2]).

execute(Req, Env) ->
    case cowboy_req:header(<<"authorization">>, Req) of
        <<"Bearer ", Token/binary>> ->
            case verify_token(Token) of
                {ok, Claims} ->
                    {ok, cowboy_req:set_meta(claims, Claims, Req), Env};
                {error, _} ->
                    {stop, reply_401(Req)}
            end;
        undefined ->
            {stop, reply_401(Req)}
    end.

reply_401(Req) ->
    cowboy_req:reply(401, #{<<"www-authenticate">> => <<"Bearer">>},
        <<"{\"error\":\"unauthorized\"}">>, Req).
%% 注册中间件链(顺序即执行顺序):日志 → CORS → 认证 → 路由 → 处理
middlewares => [request_id_mw, cors_mw, auth_middleware,
                cowboy_router, cowboy_handler].

%% 请求 ID 中间件:透传上游 ID 或生成新 ID,便于链路追踪
-module(request_id_mw).
-behaviour(cowboy_middleware).
-export([execute/2]).

execute(Req, Env) ->
    ReqId = case cowboy_req:header(<<"x-request-id">>, Req) of
                undefined -> binary:encode_hex(crypto:strong_rand_bytes(16));
                Id -> Id
            end,
    {ok, cowboy_req:set_resp_header(<<"x-request-id">>, ReqId, Req), Env}.

3.2 JSON 编解码选型

库特点适用
jsx纯 Erlang,稳定通用首选
thoas性能优于 jsx高吞吐
jiffyNIF 实现,最快极致性能
json (OTP 27+)标准库新项目
%% 编码:注意 binary 与 atom 的处理
encode_user(#{id := Id, name := Name}) ->
    jsx:encode(#{id => Id, name => Name,
                 created_at => iso8601(erlang:system_time(second))}).

%% 解码:始终用 return_maps,避免 proplist 的歧义
Decode = fun(Bin) -> jsx:decode(Bin, [return_maps]) end.

%% 安全:解码后必须做 schema 校验,不能直接信任
validate_create(#{<<"name">> := Name, <<"email">> := Email})
  when is_binary(Name), byte_size(Name) > 0 ->
    {ok, #{name => Name, email => Email}};
validate_create(_) ->
    {error, invalid_payload}.

安全提醒:jsx:decode 不会自动把 key 转成 atom(除非显式 [labels, atom]),这恰恰是安全的默认行为——绝不要把外部输入的 key 转成 atom,否则会造成 atom 表耗尽。

四、WebSocket 与 SSE 实时通信

4.1 WebSocket handler

Cowboy 的 WebSocket 基于同一套 handler 协议,通过 cowboy_websocket 升级:

-module(ws_handler).
-behaviour(cowboy_websocket).
-export([init/2, websocket_init/1, websocket_handle/2,
         websocket_info/2, terminate/3]).

init(Req, State) ->
    {cowboy_websocket, Req, State, #{idle_timeout => 60_000}}.

websocket_init(State) ->
    ok = pg:join(realtime, self()),          % 订阅业务事件
    {ok, State}.

websocket_handle({text, Msg}, State) ->
    case jsx:decode(Msg, [return_maps]) of
        #{<<"type">> := <<"ping">>} ->
            {reply, {text, <<"{\"type\":\"pong\"}">>}, State};
        _ ->
            {ok, State}
    end;
websocket_handle(_Frame, State) ->
    {ok, State}.

%% 收到 Erlang 消息 → 推送给客户端
websocket_info({event, Event}, State) ->
    {reply, {text, jsx:encode(Event)}, State};
websocket_info(_Info, State) ->
    {ok, State}.

terminate(_Reason, _Req, _State) ->
    pg:leave(realtime, self()),
    ok.
%% 业务侧广播:向所有订阅者推送
broadcast(Event) -> [Pid ! {event, Event} || Pid <- pg:get_members(realtime)], ok.

4.2 SSE:服务端推送

对于只需要单向推送的场景,SSE 比 WebSocket 更简单,且天然支持自动重连:

-module(sse_handler).
-export([init/2]).

init(Req0, State) ->
    Headers = #{<<"content-type">> => <<"text/event-stream">>,
                <<"cache-control">> => <<"no-cache">>},
    Req = cowboy_req:stream_reply(200, Headers, Req0),
    self() ! tick,
    sse_loop(Req, State).

sse_loop(Req, State) ->
    receive
        tick ->
            cowboy_req:stream_body(<<"data: ", (ts())/binary, "\n\n">>, nofin, Req),
            erlang:send_after(1000, self(), tick),
            sse_loop(Req, State)
    end.
维度WebSocketSSE
方向双向服务端 → 客户端
协议自定义帧纯 HTTP 文本
重连需自行实现浏览器自动重连
代理友好度需配置升级好
适用聊天、协作通知、进度、行情

4.3 实时通信的背压

WebSocket 推送过快会堆积在连接进程邮箱中,必须做水位保护:

websocket_info({event, #{priority := low} = Event}, State) ->
    case process_info(self(), message_queue_len) of
        {message_queue_len, N} when N > 1000 -> {ok, State};   % 丢弃低优先级
        _ -> {reply, {text, jsx:encode(Event)}, State}
    end;
websocket_info({event, Event}, State) ->
    {reply, {text, jsx:encode(Event)}, State}.

五、限流、安全与部署

5.1 令牌桶限流中间件

-module(rate_limit_mw).
-behaviour(cowboy_middleware).
-export([execute/2]).

-define(RATE, 100).        % 每秒补充 100 个令牌
-define(BURST, 200).       % 桶容量

execute(Req, Env) ->
    case allow(peer_ip(Req)) of
        true  -> {ok, Req, Env};
        false ->
            Req2 = cowboy_req:reply(429,
                #{<<"retry-after">> => <<"1">>},
                <<"{\"error\":\"rate limited\"}">>, Req),
            {stop, Req2}
    end.

allow(Ip) ->
    Now = erlang:monotonic_time(millisecond),
    Key = {bucket, Ip},
    case ets:lookup(rate_buckets, Key) of
        [{_, Tokens, LastRefill}] ->
            Refilled = min(?BURST, Tokens + (Now - LastRefill) * ?RATE div 1000),
            case Refilled >= 1 of
                true  -> ets:insert(rate_buckets, {Key, Refilled - 1, Now}), true;
                false -> false
            end;
        [] ->
            ets:insert(rate_buckets, {Key, ?BURST - 1, Now}),
            true
    end.

5.2 安全加固清单

风险防护
请求体过大read_body 设 length 上限
头部注入拒绝含 \r\n 的用户输入
CORS 滥用显式白名单 Origin,不用通配符
慢速攻击设置 request_timeout、idle_timeout
TLS 降级强制 TLS 1.2+,禁用弱套件
信息泄露生产环境不返回堆栈
%% TLS 监听器配置
{ok, _} = cowboy:start_tls(https_listener,
    [{port, 8443},
     {certfile, "/etc/ssl/server.crt"},
     {keyfile, "/etc/ssl/server.key"},
     {versions, ['tlsv1.2', 'tlsv1.3']}],
    #{env => #{dispatch => Dispatch}}).

5.3 部署与优雅停机

%% 应用停止时先摘除监听器,等待在途请求完成
prep_stop(State) -> cowboy:stop_listener(http_listener),
                    timer:sleep(5000), State.
%% release 方式部署(内嵌 ERTS),容器环境用 foreground 让容器成为 PID 1
rebar3 as prod release && _build/prod/rel/api/bin/api foreground
curl -f http://localhost:8080/api/health || exit 1   % 健康检查

连接数规划:每个连接占一个 BEAM 进程,10 万并发连接约占几百 MB 内存。同时需调整 +P(最大进程数)与文件描述符上限(ulimit -n),二者任一是瓶颈都会导致新连接被拒。

六、总结

Cowboy 用 Erlang 的进程模型把 HTTP 服务做成了「每连接一进程」的直观结构,这让并发与容错都变得自然。构建生产级 JSON API 的要点:

  • 架构:连接即进程,请求在连接进程内被 handler 处理,middleware 串起横切逻辑;
  • 路由与 handler:用 cowboy_router:compile 声明路由,init/2 按 method 分派,响应统一封装;
  • JSON:选型看吞吐需求,解码用 return_maps,绝不做 atom 转换,解码后必须 schema 校验;
  • 实时:WebSocket 适合双向,SSE 适合单向推送,两者都要做邮箱水位背压;
  • 安全:限流中间件、请求体上限、超时设置、TLS 配置、CORS 白名单缺一不可;
  • 部署:release + foreground,配合健康检查与优雅停机。

Cowboy 的 handler 本质是 OTP 进程,因此它与 https://plumephp.com/erlang-otp-framework/ 中的监督树天然契合;把 API 服务纳入应用的监督树,即可获得崩溃自愈能力。实时推送的进一步演进可参考 https://plumephp.com/elixir-phoenix-channels-realtime/ 中 Phoenix Channels 的设计思路。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. rebar3 构建与发布:Erlang 工程化的完整工具箱
  2. RabbitMQ 与消息中间件:AMQP 模型与可靠投递
  3. Erlang 数据库集成与连接池实战