引言
聊天、通知、协作编辑、在线游戏——现代应用的核心是实时。Phoenix Channels 建立在 OTP 分布式与 BEAM 多进程之上,用一行 subscribe 就能把消息推给成千上万连接。本文讲透 Channel 架构:从 WebSocket 握手、加入 Topic,到 PubSub 广播与 Presence 在线状态。
前置:/elixir-intro-phoenix/(Phoenix 基础)、/erlang-otp-framework/(进程与 pubsub)、/erlang-distributed-programming/(分布式消息)。
目录
- 1. Channel 架构全景
- 2. Socket 与连接认证
- 3. Channel 加入与生命周期
- 4. 客户端消息流:handle_in 与 push
- 5. PubSub:服务端广播
- 6. Topic 设计模式
- 7. Presence 在线状态
- 8. 性能与背压
- 9. 与 LiveView 的配合
- 10. 速查表与一句话记忆
- 延伸阅读
1. Channel 架构全景
1.1 分层结构
浏览器 WebSocket 客户端
⇅ (Phoenix JS 客户端)
Socket(连接:认证、传输管理)
⇅ (1 个连接 N 个 Channel)
Channel(业务主题:聊天室/房间/通知)
⇅ (PubSub 广播 / 分发到订阅者进程)
Erlang/Elixir 进程(每个 Channel 一个进程)
1.2 核心组件
| 组件 | 职责 |
|---|---|
| Socket | 管理 WebSocket 连接与认证 |
| Channel | 处理具体主题业务 |
| PubSub | 服务端广播分发 |
| Presence | 在线状态追踪 |
记忆:Channel 架构 = Socket(连接/认证)+ Channel(主题业务进程)+ PubSub(广播分发)+ Presence(在线状态)——浏览器通过 WebSocket 连到 Socket,加入多个 Channel。
2. Socket 与连接认证
2.1 定义 Socket
# lib/my_app_web/channels/user_socket.ex
defmodule MyAppWeb.UserSocket do
use Phoenix.Socket
channel "room:*", MyAppWeb.RoomChannel
channel "notification:*", MyAppWeb.NotificationChannel
# 连接时调用:校验认证 token
def connect(%{"token" => token}, socket, _connect_info) do
case MyApp.Auth.verify(token) do
{:ok, user} -> {:ok, assign(socket, :user, user)}
:error -> :error # 拒绝连接
end
end
def id(socket), do: "users_socket:#{socket.assigns.user.id}"
end
2.2 客户端连接
import { Socket } from "phoenix";
let socket = new Socket("/socket", { params: { token: "..." } });
socket.connect();
记忆:Socket 定义 connect 回调做认证(token 校验)、channel 宏声明路由(“room:*” 通配 Topic)、id 返回连接唯一标识——客户端 new Socket + params 带 token。
3. Channel 加入与生命周期
3.1 Channel 回调
defmodule MyAppWeb.RoomChannel do
use Phoenix.Channel
# 加入 Topic 时调用:可以拒绝
def join("room:" <> room_id, _payload, socket) do
case MyApp.Room.can_join?(socket.assigns.user, room_id) do
true ->
# 订阅进程到 PubSub Topic
{:ok, assign(socket, :room_id, room_id)}
false ->
{:error, %{reason: "unauthorized"}}
end
end
def terminate(_reason, socket) do
# 清理:通知下线等
{:ok, socket}
end
end
3.2 客户端加入
let channel = socket.channel("room:42", {});
channel.join()
.receive("ok", resp => console.log("加入成功"))
.receive("error", resp => console.log("被拒绝"));
记忆:Channel join 回调做权限校验(返回 {:ok, socket} 或 {:error, reason}),terminate 做清理;客户端 channel(topic) + join().receive 处理结果——每个 Channel 一个进程。
4. 客户端消息流:handle_in 与 push
4.1 客户端发消息
# 服务端处理客户端消息
def handle_in("new_msg", %{"body" => body}, socket) do
broadcast!(socket, "new_msg", %{body: body, user: socket.assigns.user.name})
{:noreply, socket}
end
4.2 服务端主动推
# 从外部(任务、其他 Channel)主动推给 Topic 订阅者
MyAppWeb.Endpoint.broadcast("room:42", "system", %{text: "公告"})
4.3 客户端接收
channel.on("new_msg", payload => appendMsg(payload));
记忆:handle_in 处理客户端事件、broadcast! 转发给 Topic 所有人;服务端任何地方用 Endpoint.broadcast 主动推;客户端 channel.on 接收——全双工消息流。
5. PubSub:服务端广播
5.1 PubSub 架构
Phoenix 内置 Phoenix.PubSub(基于 PG2 / 分布式注册),进程订阅 Topic,broadcast 时分发:
broadcast("room:42", event, payload)
↓
PubSub 路由(本节点广播 + 跨节点转发)
↓
所有订阅该 Topic 的 Channel 进程 → push 给各自 WebSocket
5.2 手动订阅
# 普通进程也能订阅(不一定是 Channel)
:ok = MyApp.PubSub.subscribe(MyApp.PubSub, "orders:new")
# 收到消息
receive do
{:orders:new, payload} -> handle_order(payload)
end
5.3 跨节点
多节点部署时 PubSub 自动跨节点分发(需配置 pg2 / 本地 pubsub 模式)。
记忆:PubSub = 进程订阅 + 广播分发——Channel join 时隐式订阅 Topic、broadcast! 推给所有订阅者、普通进程可手动 subscribe;多节点自动跨节点转发。
6. Topic 设计模式
6.1 常用 Topic 结构
"room:" <> id 聊天室
"user:" <> id 用户私信/通知
"order:" <> id 订单状态(电商实时)
"metrics:" <> id 监控指标流
6.2 设计原则
✓ 粒度按「谁需要收到」划分
✓ 大 Topic 拆小(room:42 比 global 好)
✓ 权限在 join 里校验,不放客户端
✓ Topic 名避免敏感信息(用 ID 不用用户名)
6.3 典型:订单状态推送
# 订单服务里(非 Channel 上下文)
MyAppWeb.Endpoint.broadcast("order:" <> order_id, "status", %{status: "paid"})
# 用户端订阅
let ch = socket.channel("order:" + order.id);
ch.join(); ch.on("status", s => updateUI(s));
记忆:Topic 设计按「谁需要收到」分粒度(room:id / user:id / order:id)、权限在 join 校验、topic 用 ID 不用敏感名——订单状态等业务事件用 Endpoint.broadcast 推给对应 Topic。
7. Presence 在线状态
7.1 追踪在线用户
Phoenix Presence 用 CRDT 同步各节点的在线列表:
# 初始化 Presence
defmodule MyAppWeb.Presence do
use Phoenix.Presence, otp_app: :my_app,
pubsub_server: MyApp.PubSub
end
# Channel join 后追踪
def handle_info({:presence_diff, diff}, socket) do
push(socket, "presence_diff", diff)
{:noreply, socket}
end
# 加入时把用户加入 Presence 列表
MyAppWeb.Presence.track(self(), "room:42",
%{user_id: user.id, name: user.name, online_at: now()})
7.2 同步展示
let presence = new Presence(channel);
presence.onSync(() => renderUsers(presence.list()));
记忆:Presence 用 CRDT 跨节点同步在线列表——track 登记、presence_diff 广播差异、客户端 Presence.onSync 同步渲染——在线状态跨节点一致无需手动管理。
8. 性能与背压
8.1 扩展模型
单节点瓶颈:单个 PubSub 进程分发
水平扩展:加节点(PubSub 自动跨节点)+ 负载均衡 WebSocket
8.2 背压与限速
# 广播太频繁会压垮客户端 —— 用节流
def handle_in("typing", payload, socket) do
# 用 :throttle 限 1 次/2s
case :ets.lookup(:typing_throttle, socket.assigns.room_id) do
[] -> :ets.insert(:typing_throttle, {room_id, now()})
broadcast!(socket, "typing", payload)
{:noreply, socket}
_ -> {:noreply, socket} # 丢弃
end
end
8.3 大连接数
- 增加 socket 长连接参数、WebSocket 心跳(
/socket已内置 ping/pong) - 用
Phoenix.Transports配置压缩 - 监控:
Phoenix.PubSub订阅数、连接数指标
记忆:性能三招——水平扩展(PubSub 自动跨节点)、广播限速(typing 节流防压垮)、WebSocket 心跳与连接监控——大连接数靠加节点与背压控制。
9. 与 LiveView 的配合
9.1 分工
| 技术 | 场景 |
|---|---|
| LiveView | 页面状态驱动的交互(表单、局部刷新) |
| Channels | 全局/跨页实时事件(聊天、通知、多端同步) |
9.2 结合模式
# LiveView 里订阅 Channel 事件,驱动 UI
def mount(_params, _session, socket) do
MyApp.PubSub.subscribe(MyApp.PubSub, "user:" <> user_id)
{:ok, socket}
end
def handle_info({:user:notification, notif}, socket) do
{:noreply, push_notification(socket, notif)}
end
记忆:LiveView 管页面状态交互、Channels 管跨页实时事件——LiveView mount 时 subscribe PubSub Topic、handle_info 里驱动 UI,两者互补。
10. 速查表与一句话记忆
| 环节 | 关键 API |
|---|---|
| 定义 Socket | use Phoenix.Socket + connect/2 |
| 路由 Topic | channel “room:*”, RoomChannel |
| 加入校验 | join/3 → {:ok, socket} / {:error, r} |
| 客户端消息 | handle_in/3 + broadcast!/3 |
| 主动推送 | Endpoint.broadcast/3 |
| 在线状态 | Presence.track + presence_diff |
| 普通订阅 | PubSub.subscribe/2 |
| 限速 | :ets 节流 |
一句话记忆:Phoenix Channels 实时通信 = Socket(WebSocket 连接 + connect 认证)+ Channel(每主题一个进程:join 校验权限、handle_in 处理客户端消息、broadcast! 转发给 Topic 所有人)+ PubSub(进程订阅/广播、跨节点自动分发)+ Presence(CRDT 同步在线列表);Topic 按「谁需要收到」设计(room:id/user:id/order:id)、权限在 join 里校验;性能靠水平扩展 + 广播限速(节流防压垮);LiveView 管页面状态交互、Channels 管跨页实时事件,两者 subscribe Topic 互补。"
延伸阅读
- /elixir-intro-phoenix/ — Phoenix 框架入门
- /elixir-liveview-advanced/ — LiveView 进阶
- /erlang-otp-framework/ — OTP 进程与 PubSub
- /erlang-distributed-programming/ — 分布式消息传递
- /erlang-mnesia-distributed/ — 分布式存储
- [[web]] — Web 全栈开发
- [[nodejs]] — Node.js 实时通信对比
- Phoenix Channels 文档
- Phoenix Presence 文档
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。