在 K8S 环境中构建基于 Python Flask 的 WebSocket 推送架构
业务上需要搭建一个能推送消息的架构,现有服务端基于 Python Flask 构建。下面记录代码实现和 K8S + HAProxy 两层负载均衡下的会话保持方案。
代码部分
服务端
服务端需要集成 socketio,参考 python-socketio 官方示例。记得把 async_mode 改成 gevent。
前端
前端测试页面,需要替换 <Server IP>:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86
| <!DOCTYPE HTML> <html> <head> <title>Flask-SocketIO Test</title> <script type="text/javascript" src="//code.jquery.com/jquery-2.1.4.min.js"></script> <script type="text/javascript" src="//cdnjs.cloudflare.com/ajax/libs/socket.io/3.0.3/socket.io.min.js"></script> <script type="text/javascript" charset="utf-8"> $(document).ready(function(){ var socket = io.connect("http://<Server IP>:5000/");
socket.on('connect', function() { socket.emit('my_event', {data: 'I\'m connected!'}); }); socket.on('disconnect', function() { $('#log').append('<br>Disconnected'); }); socket.on('nuke_response', function(msg) { $('#log').append('<br>Received: ' + msg.data); });
$('form#emit').submit(function(event) { socket.emit('my_event', {data: $('#emit_data').val()}); return false; }); $('form#broadcast').submit(function(event) { socket.emit('my_broadcast_event', {data: $('#broadcast_data').val()}); return false; }); $('form#join').submit(function(event) { socket.emit('join', {room: $('#join_room').val()}); return false; }); $('form#leave').submit(function(event) { socket.emit('leave', {room: $('#leave_room').val()}); return false; }); $('form#send_room').submit(function(event) { socket.emit('my_room_event', {room: $('#room_name').val(), data: $('#room_data').val()}); return false; }); $('form#close').submit(function(event) { socket.emit('close_room', {room: $('#close_room').val()}); return false; }); $('form#disconnect').submit(function(event) { socket.emit('disconnect_request'); return false; }); }); </script> </head> <body> <h1>Flask-SocketIO Test</h1> <h2>Send:</h2> <form id="emit" method="POST" action='#'> <input type="text" name="emit_data" id="emit_data" placeholder="Message"> <input type="submit" value="Echo"> </form> <form id="broadcast" method="POST" action='#'> <input type="text" name="broadcast_data" id="broadcast_data" placeholder="Message"> <input type="submit" value="Broadcast"> </form> <form id="join" method="POST" action='#'> <input type="text" name="join_room" id="join_room" placeholder="Room Name"> <input type="submit" value="Join Room"> </form> <form id="leave" method="POST" action='#'> <input type="text" name="leave_room" id="leave_room" placeholder="Room Name"> <input type="submit" value="Leave Room"> </form> <form id="send_room" method="POST" action='#'> <input type="text" name="room_name" id="room_name" placeholder="Room Name"> <input type="text" name="room_data" id="room_data" placeholder="Message"> <input type="submit" value="Send to Room"> </form> <form id="close" method="POST" action="#"> <input type="text" name="close_room" id="close_room" placeholder="Room Name"> <input type="submit" value="Close Room"> </form> <form id="disconnect" method="POST" action="#"> <input type="submit" value="Disconnect"> </form> <h2>Receive:</h2> <div><p id="log"></p></div> </body> </html>
|
这样就完成了正常环境下 WebSocket 的代码部分。
架构部分
架构图

几条重点:
- 系统发布在 K8S 中,一个 Application 的 Deployment 包含三个 Pod 和一个 Service。
- K8S Service 用 NodePort 暴露服务。
- K8S 外面用 HAProxy 代理,域名解析到 HAProxy 所在虚拟机的 IP。
- HAProxy 和 K8S Service 形成两层 Load Balancing。
会话保持
目标是让客户端 A、B、C 访问到某个 Pod 后,以后就一直绑定到这个 Pod。需要改三处配置:
| 层级 |
配置 |
作用 |
| HAProxy |
balance source |
按客户端 IP 选择后端 |
| K8S Service |
sessionAffinity: ClientIP |
Service 层按客户端 IP 转发 |
| K8S Service |
externalTrafficPolicy: Local |
保留真实客户端 IP(否则所有请求都会被当成来自 HAProxy 的 IP) |
消息队列
正常情况下,客户端 A 和 B 连到第一个 Pod,客户端 C 连到第三个 Pod。如果事件发生在第二个 Pod 或 Job Pod 上,无法直接推送消息给所有客户端。
因此需要一个消息总线(Redis / MQ / kube-event 均可):三个 API Pod 侦听消息队列的某个事件,需要发送给客户的消息先写入队列,队列再转发给三个 API Pod,各 Pod 收到提醒后推送给自己连接的客户端,完成整体回路。