做后端开发这些年,socket 这个词几乎天天都能碰到,但真正让我把socket和“发布与订阅”(Pub/Sub)结合起来做项目,还是在一次消息推送需求里被逼出来的。当时要做一个多端实时通知系统,HTTP 轮询太重,单靠数据库发短信又贵,最后落地方案就是:底层用 socket 维持长连接,上层按主题做发布订阅。这个组合听起来简单,实际落地时会遇到协议设计、粘包、心跳、订阅关系维护一堆问题。这篇文章就把我从零开始做 socket 发布订阅的完整思路和踩坑记录整理出来,适合刚接触网络编程的后端开发、物联网开发者,以及想在 Flask、Spring Boot 项目里做实时消息推送的朋友。
1. 先搞清楚 socket 发布订阅到底解决什么问题
1.1 发布订阅不是简单的“发消息”
很多新手会把发布订阅和“两个人聊天”搞混。TCP socket 本身是点对点的通道,A 和 B 建立连接后,A 发什么 B 就收什么,这是“单播”。但真实业务里往往是一对多:一个用户发布了失物招领信息,系统要推给所有关注“校园卡”这个标签的人;一个传感器上报了温度,平台要同步推给大屏、手机 App、报警服务。如果每个接收方都单独建一条连接去发,代码会膨胀得很厉害,连接数也撑不住。
发布订阅模式的出现就是为了解决这种“多对多”的解耦问题。它引入了“主题”的概念:发布者把消息发到某个主题上,订阅者提前声明“我只对这个主题感兴趣”,消息中间层负责把主题上的消息扇出(fanout)给所有订阅者。发布者根本不用关心谁在听,订阅者也不用关心消息从哪来,两边只和“主题”打交道。
1.2 socket 在发布订阅里扮演什么角色
发布订阅的“消息中间层”可以是 Redis、RabbitMQ 这类独立中间件,但也有很多场景需要自己用 socket 实现。socket 在这里承担的是“传输通道”职责:它负责把发布者发的消息从进程 A 的网卡搬到进程 B 的缓冲区,再交给订阅者的业务逻辑。用 socket 做底层的好处是灵活,你可以自定义协议、控制心跳、控制消息格式,不依赖重型中间件;坏处是很多细节要自己补,比如断线重连、消息确认、网络字节序。
打个生活化的比方:socket 像一条条管道,发布订阅像管道上装的“分拣器”。没有分拣器时,每根管道只能点对点送;装上分拣器后,一个人往管道里丢包裹,分拣器按标签把包裹放进不同格口,每个格口对应一个订阅者。socket 保证管道不堵不漏,分拣器保证包裹送对人。
1.3 典型应用场景一览
- 实时通知:用户关注某件招领物品,物品状态更新时立即推送。
- 物联网数据上报:设备通过 MQTT over TCP 上报状态,平台按设备主题发布给多个下游系统。
- 聊天室:多个用户订阅同一个房间主题,消息广播给房间内所有人。
- 股票行情:行情源发布价格,所有订阅了该股票代码的客户端实时收到。
这些场景的共同特征是:实时性要求高、客户端状态多样、消息需要按兴趣过滤。用 HTTP 轮询能实现但不优雅,用 socket 长连接加发布订阅才是业内常见的正解。
2. 核心技术拆解:socket 与发布订阅模式的连接点
2.1 发布订阅的四个核心角色
不管用 TCP、WebSocket 还是 MQTT,发布订阅都有四个固定角色:
- 发布者(Publisher):只负责往主题发消息。
- 订阅者(Subscriber):只负责表达“我要订阅什么”,并接收消息。
- 主题(Topic):消息的分类标签,也可以是一串带层级的关键词。
- 消息代理(Broker/Server):维护主题与订阅者的映射关系,负责转发。
自研 socket 发布订阅时,你的核心工作就是实现第四个角色:一个常驻的 socket 服务端,它既要接受发布者的连接,也要接受订阅者的连接,还要维护一张“主题 -> 订阅者连接集合”的关系表。这张表是整个系统的灵魂。
2.2 为什么不用 HTTP 轮询
- 轮询延迟高:即使间隔 1 秒,用户感知也有至少 1 秒延迟,而且服务端无法主动推送。
- 资源浪费:大量无意义请求会压垮轻量化服务器。
- 代码割裂:查询接口和推送接口是两个体系,业务逻辑难统一。
socket 长连接建立后,服务端可以随时主动往客户端写数据,这才叫“推送”。发布订阅模式最大的收益就是服务端由“被动应答”变成“主动分发”,而这必须依赖长连接。
2.3 三种落地方式怎么选
| 方案 | 底层协议 | 优点 | 缺点 | 适合场景 |
|---|---|---|---|---|
| 自研 TCP + 自定义协议 | TCP | 可控性最强,性能高 | 要自己处理粘包、心跳、序列化 | 嵌入式、游戏、内部系统 |
| WebSocket 发布订阅 | WebSocket(基于 TCP) | 浏览器原生支持,穿透防火墙 | 需要处理握手和帧解析,有少量额外开销 | 网页端实时推送 |
| MQTT | TCP/TLS | 协议成熟,QoS 可靠,生态完善 | 需要部署 Broker(如 EMQX、Mosquitto) | 物联网、移动端弱网环境 |
自研 TCP 方案看起来“最底层”,但发布订阅的核心逻辑完全不依赖协议类型:不管底层是 TCP socket 还是 WebSocket,上层都要维护订阅关系、按主题分发。所以我的做法是先把自研 TCP 方案跑通,再把同一套逻辑移植到 WebSocket 上,理解会非常透彻。
3. 手写一个基于 TCP socket 的发布订阅 Demo
3.1 先设计消息协议:避免“后面补随机数”的坑
自研 TCP 发布订阅,第一件要事是定义消息帧。很多人在网上查资料时会看到“为什么 socket 接收到奇数字节,后面会补一个随机数”这种问题,其实那不是随机数,而是 TCP 是字节流协议,没有消息边界。如果发送方一次 send 了 5 个字节,接收方可能一次 recv 到 3 个字节,另一次 recv 到 2 个字节,如果发送的字段正好是奇数长度,接收方硬按固定结构去切分,就会把下一个消息的头部当成“补的随机数”解析。
解决办法是设计一个带有“长度字段”的协议帧。我用 Python 的struct打包,帧结构如下:
| 2字节魔数 | 1字节类型 | 2字节主题长度 | 主题内容 | 4字节载荷长度 | 载荷内容 |其中:
- 魔数固定为
0x5050,用来快速校验是否为合法帧。 - 类型:1 表示订阅,2 表示取消订阅,3 表示发布,4 表示服务器推送的消息,5 表示心跳。
- 主题长度和载荷长度都按网络字节序(大端)编码,避免不同机器解析出错。
接收端必须循环读取,直到读满一个完整帧再加处理,这就是经典的“拆包”。核心思路是:先读固定长度的头部,解析出主题长度和载荷长度,再按这个长度读剩余字节,读不满就继续等下一次 recv。这样奇数字节、补随机数的问题就不会出现。
3.2 服务端:维护订阅关系的 Broker
服务端代码的核心是维护subscriptions字典,键是主题字符串,值是该主题下所有订阅者的 socket 连接对象列表。我用 Python 的threading为每个客户端连接开一个线程,这样发布和订阅可以并行处理。代码逻辑如下:
import socket import struct import threading SUBSCRIPTIONS = {} LOCK = threading.Lock() HEADER = struct.Struct(">HBH") # 魔数+类型+主题长度 def recv_exact(conn, size): data = b"" while len(data) < size: chunk = conn.recv(size - len(data)) if not chunk: raise ConnectionError("连接已断开") data += chunk return data def recv_frame(conn): header = recv_exact(conn, HEADER.size) magic, msg_type, topic_len = HEADER.unpack(header) if magic != 0x5050: raise ValueError("魔数校验失败") topic = recv_exact(conn, topic_len).decode("utf-8") payload_len = struct.unpack(">I", recv_exact(conn, 4))[0] payload = recv_exact(conn, payload_len).decode("utf-8") return msg_type, topic, payload def handle_conn(conn): while True: try: msg_type, topic, payload = recv_frame(conn) except (ConnectionError, ValueError): with LOCK: for t, clients in list(SUBSCRIPTIONS.items()): if conn in clients: clients.remove(conn) if not clients: del SUBSCRIPTIONS[t] conn.close() return if msg_type == 1: # 订阅 with LOCK: SUBSCRIPTIONS.setdefault(topic, []).append(conn) elif msg_type == 3: # 发布 with LOCK: for sub_conn in SUBSCRIPTIONS.get(topic, []): try: send_push(sub_conn, topic, payload) except Exception: pass def send_push(conn, topic, payload): topic_b = topic.encode("utf-8") payload_b = payload.encode("utf-8") frame = HEADER.pack(0x5050, 4, len(topic_b)) + topic_b frame += struct.pack(">I", len(payload_b)) + payload_b conn.sendall(frame) def start_broker(host="0.0.0.0", port=9000): srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) srv.bind((host, port)) srv.listen(128) while True: conn, _ = srv.accept() threading.Thread(target=handle_conn, args=(conn,), daemon=True).start() if __name__ == "__main__": start_broker()这段代码虽然简化,但已经把订阅、取消订阅、发布、推送的核心流程都实现了。注意我用了LOCK保证多线程下SUBSCRIPTIONS的读写安全,否则高并发时字典会被同时修改,轻则丢消息,重则程序崩溃。
3.3 订阅者客户端与发布者客户端
订阅者客户端建立连接后,发一个订阅帧,然后循环等待服务器推送:
import socket import struct import time HEADER = struct.Struct(">HBH") def subscribe(host, port, topic): conn = socket.create_connection((host, port)) topic_b = topic.encode("utf-8") frame = HEADER.pack(0x5050, 1, len(topic_b)) + topic_b + struct.pack(">I", 0) conn.sendall(frame) print(f"已订阅主题: {topic}") while True: header = recv_exact(conn, HEADER.size) magic, msg_type, topic_len = HEADER.unpack(header) recv_topic = recv_exact(conn, topic_len).decode("utf-8") payload_len = struct.unpack(">I", recv_exact(conn, 4))[0] payload = recv_exact(conn, payload_len).decode("utf-8") if msg_type == 4: print(f"[{time.strftime('%H:%M:%S')}] 收到主题 '{recv_topic}' 的消息: {payload}")发布者客户端更简单,只需要按类型 3 发送一帧即可。这里的关键点是:发布者并不需要与某个订阅者建立直连,它只要连接到 broker,把消息扔给 broker,broker 再去查订阅表转发。这样就实现了发布者和订阅者的完全解耦。
3.4 亲手跑一遍 Demo 的步骤
- 启动 broker:
python broker.py - 开两个订阅终端:分别订阅主题
lost和found - 开一个发布终端:向主题
lost发布一条“校园卡丢失” - 观察只有订阅了
lost的终端收到消息
这个过程会把发布订阅的核心链路完整验证一遍。我建议初学者不要直接上框架,先把这种裸 socket 版本跑通,后面即使换成 Netty、Tornado 或 Spring WebSocket,理解成本都会低很多。
4. WebSocket 与 MQTT:现实中更常用的发布订阅通道
4.1 浏览器端首选 WebSocket
裸 TCP 无法直接在浏览器里使用,所以网页端做实时推送普遍选择 WebSocket。WebSocket 本质上是基于 TCP 的上层协议,完成握手后,服务端可以随时推送消息给浏览器。发布订阅的模型并没有变,变的只是消息封装格式和多了一个握手过程。
我在 Spring Boot 项目里集成过 WebSocket,网上最多的坑集中在yml配置上。很多人问“Spring Boot 集成 WebSocket 的 yml 配置怎么写”,其实 Spring Boot 的 WebSocket 一般不通过 yml 配置端口和路径,而是通过配置类注册端点。常用的配置片段如下:
server: port: 8080 spring: application: name: ws-demo核心配置在 Java 代码里:
@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new TopicSocketHandler(), "/ws/topic") .setAllowedOrigins("*"); } }真正要调整的是setAllowedOrigins,如果前端跨域,必须要允许对应来源,否则浏览器握手会被拦。这类“socket 有跨域吗”的问题,本质就是 WebSocket 的跨域校验规则和 HTTP 不完全一样,需要单独配置。
4.2 MQTT:物联网发布订阅的事实标准
如果项目涉及物联网设备上报或者弱网环境的消息推送,直接用 MQTT 比自研协议划算得多。MQTT 基于 TCP,协议体量很小,固定头部最少只要 2 字节,非常适合传感器、嵌入式设备。它最大的优势是内置了 QoS(服务质量)等级:QoS 0 最多发一次、QoS 1 至少发一次、QoS 2 恰好发一次。自研协议要做到可靠投递非常复杂,MQTT 直接帮你解决了。
订阅规则也很有意思,支持通配符:订阅sensor/+/temperature可以匹配sensor/room1/temperature和sensor/room2/temperature,这种带层级的话题设计非常适合设备分类管理。部署时我常用 EMQX 或 Mosquitto 做 Broker,客户端用 paho-mqtt,发布订阅代码几行就写完。
4.3 自研还是用现成 Broker,我的判断标准
| 判断维度 | 自研 TCP/WebSocket 方案 | 直接使用 MQTT Broker |
|---|---|---|
| 协议定制需求 | 高 | 低 |
| 可靠投递保障 | 需要自己实现 ACK 和重发 | 内置 QoS |
| 部署复杂度 | 随代码增长而上升 | 只要部署一个服务 |
| 调试成本 | 较高,需要抓包分析 | 有现成客户端工具 |
| 适用规模 | 百级到千级连接 | 万级到百万级连接 |
如果只是做个校园失物招领平台、内部通知系统,自研方案完全够用;如果要做量产设备接入、几十万连接的地域级平台,别折腾自研协议,老老实实上 MQTT。
5. 发布订阅中的消息丢失、超时与粘包问题排查
5.1 经典的socket read timed out
不少人会在日志里看到java.sql.SQLException: IO 错误: socket read timed out或create socket connection failure。这里有个误区:这个socket不是网络推流的 socket,而是 JDBC 连接 MySQL 时报的错。它的含义是:客户端已经发起了数据库查询,但服务端在规定时间内没有返回数据。排查步骤通常是这样:
- 先确认 MySQL 是否有慢查询,
show processlist看是否有长时间卡住的会话。 - 检查连接池是否被占满,
druid或hikari的最大连接数太小,会导致后续连接排队。 - 调整 JDBC 驱动的
socketTimeout参数,比如socketTimeout=60000,但不要设成无限大,否则连接真死了你也发现不了。 - 查看网络层,确认应用服务器和数据库之间是否有防火墙丢包。
这种问题经常被误判为“发布订阅消息超时”,其实和消息推送无关,区分的关键是看报错发生在哪一层:是业务线程拿着数据库连接时超时,还是 socket 收发数据时超时。
5.2 粘包、半包和“补随机数”的真相
回到开头的“收到奇数字节,后面补一个随机数”问题。TCP 是字节流,没有天然的消息边界,所以接收方不能假设“一次 recv 就是一条完整消息”。如果发送方的自定义消息是“奇数长度 + 不固定长度”,接收方很容易把下一条消息的头部当成当前消息缺失的字节来解析。
不要想着靠“补随机数”去处理,正确做法就是我前面第 3 节里的方案:在帧头固定长度字段,接收方先读固定头,再按长度读剩余部分,读不全就继续循环 recv。所有真正生产可用的 socket 协议,包括 HTTP、MQTT、Redis 协议,都遵循这个“长度先行”的原则。
5.3 发布订阅消息丢失的典型场景
- 订阅者宕机时消息被丢弃:自研方案需要扩展 ACK 机制,订阅者处理完发送确认,Broker 未收到确认就重发。
- 发送缓冲区塞满:发布消息过快,消费端处理不过来,
sendall会阻塞甚至抛异常。要对每条发送加上超时,并对堆积做限流或丢弃策略。 - 心跳失效导致假死连接:长时间没数据的 TCP 连接可能被中间设备掐断,服务端却还认为它活着。解决办法是每 15 秒发一个心跳帧,N 次没收到就断开并清理订阅表。
我在真实项目里还遇到过一种情况:局域网内一切正常,跨公网就频繁断线。原因是公网 NAT 设备会回收长时间空闲的映射,所以心跳既是为了保活,也是为了维持 NAT 映射。
6. 一个真实应用场景:校园失物招领平台的实时匹配推送
6.1 场景拆解:发布与订阅天然契合
这个场景我很喜欢,因为它特别贴合“发布与订阅”的语义。用户在平台上发布“丢失校园卡”或“拾到钥匙”,从发布订阅的角度看就是往主题lost或found发布消息。而每个用户登录后,平台根据他关注的关键词,自动帮他订阅了几个主题,比如“校园卡”“身份证”“耳机”。当新发布的失物信息和某个订阅者关心的关键词匹配时,系统就该通过 socket 长连接把推荐结果实时推过去。
整个平台拆成三层:
- 网页端:用户提交失物/招领信息,用 Flask 渲染。
- 匹配层:基于中文关键词相似度计算,把新信息和历史信息做匹配。
- 推送层:把匹配结果通过 WebSocket 推给在线用户。
6.2 轻量级关键词相似度匹配
中文没有天然空格分词,做完全语义匹配非常重。校园失物招领这种轻量场景,根本不需要上大模型,用关键词集合 + 字符重叠度就够用。我常用一个简化的 Jaccard 相似度:把中文文本按单字切分,再取交集和并集的比值。
def text_similarity(a, b): chars_a = set(a) chars_b = set(b) if not chars_a or not chars_b: return 0 same = len(chars_a & chars_b) total = len(chars_a | chars_b) return same / total def match_lost_found(new_text, candidates): results = [] for item in candidates: score = text_similarity(new_text, item["description"]) if score >= 0.35: results.append((item, score)) results.sort(key=lambda x: x[1], reverse=True) return results[:5]单字级相似度的优点是实现简单、不依赖外部分词库;缺点是对“校园卡”和“学生卡”这种同义词判断不准。如果想提升精度,可以把“校园卡”“身份证”“银行卡”这几个高频词做成同义词词典,先把文本做归一化替换,再计算相似度。这个优化对准确率的提升非常明显。
6.3 把发布订阅机制接进平台
技术选型上,网页端用 Flask 做业务接口,WebSocket 用 Flask-Sock 扩展,消息本体直接推给浏览器。核心流程:
- 用户
POST /api/publish提交失物信息。 - Flask 将信息写入 SQLite,并触发相似度匹配。
- 匹配完成后,平台把结果广播到用户 WebSocket 连接上。
- 前端收到推送后渲染推荐卡片,用户点击即可看到详细匹配。
这一步实现的关键是维护“用户 -> WebSocket 连接”的映射。用户登录后,WebSocket 握手时带上用户 ID,服务端把它存到线程安全的字典里。当有新失物发布时,服务端根据匹配结果找到目标用户并调send()推送,这就完成了整个发布订阅闭环。
6.4 无效信息过滤与匹配精度优化
做这个平台时,最烦的不是匹配不上,而是匹配了一堆无意义结果。比如发布“求帮忙看看有没有好心人捡到我的校园卡”,这句话虽然包含“校园卡”,但夹杂了太多无关字符,导致相似度普遍偏高。
我的优化思路是:
- 提取关键词时,去掉“求帮忙”“有没有”“好心人”“谢谢”这类语气词和常见动词。
- 只保留名词关键词列表:校园卡、钱包、钥匙、耳机、身份证、学生证、眼镜、书本。
- 匹配时用关键词集合的 Jaccard 相似度,而不是整句字符相似度。
这样匹配精度大幅提高,无效信息也被过滤掉了。对轻量化平台来说,不需要复杂的 NLP,规则 + 关键词表就能解决大部分问题。
7. 实操心得与避坑清单
7.1 我反复踩过的几个坑
- 没有给 socket 设置超时,结果连接一卡就永远卡住。正确做法是
socket.settimeout(10)或用select做超时轮询。 - 多线程共享订阅字典不加锁,并发一高就抛
RuntimeError: dictionary changed size during iteration。 - WebSocket 的 Origin 校验不过,前端连不上,排查了半天发现是路径配成了
ws://host/ws,而端点注册在/ws/topic。 - 把“发布订阅”和“消息队列”混为一谈。发布订阅强调按主题扇出,消费者之间是竞争关系还是广播关系,一定要先想清楚。这里可以是广播,也可以是分组消费,但很多自研代码默认广播,导致业务方想只让一个消费者处理时反而没法做。
7.2 快速排查速查表
| 症状 | 可能原因 | 解决方向 |
|---|---|---|
| 客户端收到“补的随机数” | 未按长度字段拆包 | 设计固定头部+长度字段,循环 recv |
| 连接频繁超时 | 网络中间设备回收空闲连接 | 增加心跳包,超时重连 |
| 发布消息丢失 | 订阅表未加锁或订阅者未注册 | 检查订阅关系,增加 ACK 机制 |
| 数据库报 socket read timed out | 慢 SQL、连接池耗尽 | 调大 socketTimeout,优化 SQL |
| MySQL 报 /tmp/mysql.sock 错误 | 客户端通过 Unix socket 找不到 MySQL | 改为 TCP 连接,或指定 socket 文件路径 |
| 浏览器 WebSocket 连不上 | Origins 校验失败 | 配置 allowedOrigins |
7.3 如何在项目里平滑落地
如果业务系统已经跑起来了,不想推倒重来,我建议分三步走:
- 先把“发布”和“订阅”抽象成两个接口:
publish(topic, payload)和subscribe(topic, handler)。 - 实现先用内存字典顶住,验证业务闭环。
- 再把底层换成 WebSocket 或 MQTT,上层业务代码完全不动。
这样做的收益在于:发布订阅的核心价值在业务层,底层传输通道只是实现细节。只要接口设计对了,从 TCP 换到 MQTT 只是换一个适配器的问题,而不是重新写一遍业务逻辑。
个人经验是,做一个实时推送系统,最开始一定要先确认消息是“最多一次”还是“至少一次”投递。校园失物招领这种场景丢了消息可以下次刷新再看到,问题不大;但如果是报警系统、交易通知,就必须有确认重发机制。不要等线上丢消息了才补,设计协议的第一天就把 QoS 等级定下来,后面能省很多事。
最后再分享一个小技巧:调试 socket 发布订阅时,别只依赖打印日志。写一个简单的--debug模式,把每个客户端连接、订阅的主题、收发的消息类型、时间戳都打进本地文件。发现问题时回放日志,比对着屏幕猜参数快得多。这个习惯我从校园失物招领平台开始养起,后来做物联网消息网关时帮了大忙。