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 请求生命周期
- Ranch 接受 TCP 连接,创建连接进程;
- 连接进程解析请求行与头部,构造
Req对象; - 依次执行配置的 middleware(如
cowboy_router、cowboy_handler); - 路由匹配到 handler 模块,调用
Handler:init(Req, Opts); - 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 | 高吞吐 |
jiffy | NIF 实现,最快 | 极致性能 |
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.
| 维度 | WebSocket | SSE |
|---|---|---|
| 方向 | 双向 | 服务端 → 客户端 |
| 协议 | 自定义帧 | 纯 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 的设计思路。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。