You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Tornado WebSockets结合MySQL与Redis的业务实现咨询

Tornado WebSocket Server with Redis/MySQL Data Fetch & Token-Based Broadcast Distribution

Got it, let's walk through building this Tornado WebSocket server that handles three key jobs: managing client connections, fetching data from Redis/MySQL on request, and routing broadcast messages to clients who've subscribed to specific tokens. Here's a practical, step-by-step implementation with code examples you can adapt:

1. Core WebSocket Handler Setup

First, we'll create a WebSocket handler to manage client connections, process incoming messages (data requests and subscription commands), and maintain a map of clients to their subscribed tokens.

import tornado.web
import tornado.websocket
import tornado.ioloop
import redis
import pymysql
import json
import threading
from tornado.options import define, options

# Thread-safe client-subscription mapping (use lock to prevent race conditions)
client_subscriptions = {}
sub_lock = threading.Lock()

class WebSocketHandler(tornado.websocket.WebSocketHandler):
    def open(self):
        print(f"New client connected: {self.request.remote_ip}")
        # Initialize subscription set for this client
        with sub_lock:
            client_subscriptions[self] = set()

    def on_close(self):
        print(f"Client disconnected: {self.request.remote_ip}")
        # Clean up subscription record
        with sub_lock:
            if self in client_subscriptions:
                del client_subscriptions[self]

    def on_message(self, message):
        try:
            msg_data = json.loads(message)
            msg_type = msg_data.get("type")

            if msg_type == "fetch":
                # Handle data fetch request from client
                data_key = msg_data.get("key")
                result = self._fetch_from_db(data_key)
                self.write_message(json.dumps({
                    "type": "fetch_response",
                    "key": data_key,
                    "data": result
                }))
            elif msg_type == "subscribe":
                # Add token to client's subscription list
                target_token = msg_data.get("token")
                with sub_lock:
                    client_subscriptions[self].add(target_token)
                self.write_message(json.dumps({
                    "type": "subscribe_ack",
                    "token": target_token,
                    "status": "success"
                }))
            elif msg_type == "unsubscribe":
                # Remove token from client's subscription list
                target_token = msg_data.get("token")
                with sub_lock:
                    if target_token in client_subscriptions[self]:
                        client_subscriptions[self].remove(target_token)
                self.write_message(json.dumps({
                    "type": "unsubscribe_ack",
                    "token": target_token,
                    "status": "success"
                }))
            else:
                self.write_message(json.dumps({
                    "type": "error",
                    "message": "Unknown message type"
                }))
        except json.JSONDecodeError:
            self.write_message(json.dumps({
                "type": "error",
                "message": "Invalid JSON format"
            }))
        except Exception as e:
            self.write_message(json.dumps({
                "type": "error",
                "message": f"Server error: {str(e)}"
            }))

    def _fetch_from_db(self, data_key):
        # Priority: Redis cache first, fall back to MySQL
        r = redis.Redis(host="localhost", port=6379, db=0, decode_responses=True)
        redis_result = r.get(data_key)
        
        if redis_result:
            return redis_result
        
        # Fetch from MySQL if not in Redis
        db_conn = pymysql.connect(
            host="localhost",
            user="your_db_user",
            password="your_db_pass",
            database="your_db_name",
            cursorclass=pymysql.cursors.DictCursor
        )
        try:
            with db_conn.cursor() as cursor:
                cursor.execute("SELECT data_value FROM your_table WHERE data_key = %s", (data_key,))
                mysql_result = cursor.fetchone()
                if mysql_result:
                    # Cache result in Redis for next time
                    r.setex(data_key, 3600, mysql_result["data_value"])
                    return mysql_result["data_value"]
                return "Data not found"
        finally:
            db_conn.close()

2. Broadcast Message Listener

Next, we need to listen for incoming broadcast messages (we'll use Redis Pub/Sub here since it's lightweight and widely used for cross-service broadcasting). When a broadcast comes in, we'll route it to all clients subscribed to the token in the message.

def _broadcast_listener():
    r = redis.Redis(host="localhost", port=6379, db=0, decode_responses=True)
    pubsub = r.pubsub()
    pubsub.subscribe("tornado_broadcast_channel")  # Subscribe to your broadcast channel

    for msg in pubsub.listen():
        if msg["type"] == "message":
            try:
                broadcast_payload = json.loads(msg["data"])
                target_token = broadcast_payload.get("token")
                message_data = broadcast_payload.get("data")

                # Send to all subscribed clients
                with sub_lock:
                    for client, subscribed_tokens in client_subscriptions.items():
                        if target_token in subscribed_tokens:
                            client.write_message(json.dumps({
                                "type": "broadcast",
                                "token": target_token,
                                "data": message_data
                            }))
            except json.JSONDecodeError:
                print("Invalid JSON in broadcast message")
            except Exception as e:
                print(f"Error processing broadcast: {str(e)}")

3. Launch the Server

Finally, we'll start the Tornado server and spin up the broadcast listener in a background thread (so it doesn't block the main event loop).

def main():
    define("port", default=8888, help="Server port", type=int)
    tornado.options.parse_command_line()

    # Start broadcast listener in daemon thread
    broadcast_thread = threading.Thread(target=_broadcast_listener)
    broadcast_thread.daemon = True
    broadcast_thread.start()

    app = tornado.web.Application([
        (r"/ws", WebSocketHandler),
    ])
    app.listen(options.port)
    print(f"Tornado WebSocket server running on ws://localhost:{options.port}/ws")
    tornado.ioloop.IOLoop.current().start()

if __name__ == "__main__":
    main()

Key Notes for Production

  • Connection Pools: Avoid creating new Redis/MySQL connections on every request. Use connection pools (e.g., redis.ConnectionPool, pymysql.pool) to improve performance.
  • Thread Safety: We used a lock for client_subscriptions to prevent race conditions between the WebSocket handler and broadcast thread—this is critical for production.
  • Message Validation: Add stricter validation for incoming messages (e.g., check required fields) to avoid unexpected errors.
  • Error Recovery: Implement reconnection logic for Redis/MySQL if connections drop, and add logging for debugging.

内容的提问来源于stack exchange,提问作者aja

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 08:47:18