Phoenix Channels 实时通信实战:WebSocket 与 PubSub 深入

Phoenix Channels 实时通信实战:Channel 与 Socket 架构、WebSocket 连接与认证、Channel 加入与消息流、PubSub 广播与订阅、Topic 设计模式(用户/房间/资源隔离)、Presence 在线状态、Channel 性能调优与背压、与 LiveView 的配合、生产部署注意事项。

引言

聊天、通知、协作编辑、在线游戏——现代应用的核心是实时。Phoenix Channels 建立在 OTP 分布式与 BEAM 多进程之上,用一行 subscribe 就能把消息推给成千上万连接。本文讲透 Channel 架构:从 WebSocket 握手、加入 Topic,到 PubSub 广播与 Presence 在线状态。

前置:/elixir-intro-phoenix/(Phoenix 基础)、/erlang-otp-framework/(进程与 pubsub)、/erlang-distributed-programming/(分布式消息)。


目录


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
定义 Socketuse Phoenix.Socket + connect/2
路由 Topicchannel “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 文档

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. 自定义 OTP Behaviour 实战:Callback 规范与行为封装
  2. Mix 工具链与 Elixir 工程化实战
  3. LiveView 进阶实战:状态管理、并发与性能优化