
做后端开发这些年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浏览器原生支持穿透防火墙需要处理握手和帧解析有少量额外开销网页端实时推送MQTTTCP/TLS协议成熟QoS 可靠生态完善需要部署 Broker如 EMQX、Mosquitto物联网、移动端弱网环境自研 TCP 方案看起来“最底层”但发布订阅的核心逻辑完全不依赖协议类型不管底层是 TCP socket 还是 WebSocket上层都要维护订阅关系、按主题分发。所以我的做法是先把自研 TCP 方案跑通再把同一套逻辑移植到 WebSocket 上理解会非常透彻。3. 手写一个基于 TCP socket 的发布订阅 Demo3.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(host0.0.0.0, port9000): 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(targethandle_conn, args(conn,), daemonTrue).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把消息扔给 brokerbroker 再去查订阅表转发。这样就实现了发布者和订阅者的完全解耦。3.4 亲手跑一遍 Demo 的步骤启动 brokerpython 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参数比如socketTimeout60000但不要设成无限大否则连接真死了你也发现不了。查看网络层确认应用服务器和数据库之间是否有防火墙丢包。这种问题经常被误判为“发布订阅消息超时”其实和消息推送无关区分的关键是看报错发生在哪一层是业务线程拿着数据库连接时超时还是 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(keylambda x: x[1], reverseTrue) 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优化 SQLMySQL 报 /tmp/mysql.sock 错误客户端通过 Unix socket 找不到 MySQL改为 TCP 连接或指定 socket 文件路径浏览器 WebSocket 连不上Origins 校验失败配置 allowedOrigins7.3 如何在项目里平滑落地如果业务系统已经跑起来了不想推倒重来我建议分三步走先把“发布”和“订阅”抽象成两个接口publish(topic, payload)和subscribe(topic, handler)。实现先用内存字典顶住验证业务闭环。再把底层换成 WebSocket 或 MQTT上层业务代码完全不动。这样做的收益在于发布订阅的核心价值在业务层底层传输通道只是实现细节。只要接口设计对了从 TCP 换到 MQTT 只是换一个适配器的问题而不是重新写一遍业务逻辑。个人经验是做一个实时推送系统最开始一定要先确认消息是“最多一次”还是“至少一次”投递。校园失物招领这种场景丢了消息可以下次刷新再看到问题不大但如果是报警系统、交易通知就必须有确认重发机制。不要等线上丢消息了才补设计协议的第一天就把 QoS 等级定下来后面能省很多事。最后再分享一个小技巧调试 socket 发布订阅时别只依赖打印日志。写一个简单的--debug模式把每个客户端连接、订阅的主题、收发的消息类型、时间戳都打进本地文件。发现问题时回放日志比对着屏幕猜参数快得多。这个习惯我从校园失物招领平台开始养起后来做物联网消息网关时帮了大忙。