Python 2.7 WebSocket客户端消息监听及跨代码调用方法咨询
在Python 2.7中让主程序监听WebSocket消息的实用方案
针对你要把独立WebSocket工具收到的消息应用到主程序的需求,我整理了几个适配Python 2.7的可行方案,按推荐程度排序:
1. 线程安全队列(最推荐)
使用Python 2.7的Queue模块实现线程间消息传递是最稳妥的方式——队列本身是线程安全的,不需要手动处理锁,能可靠地在WebSocket线程和主程序线程之间传递消息。
WebSocket工具代码(可封装为模块)
try: import thread except ImportError: import _thread as thread import time import websocket import Queue # Python 2.7专属队列模块 def on_message(ws, message): # 将收到的消息放入绑定的队列 ws.message_queue.put(message) def on_error(ws, error): print(error) def on_close(ws): print("### closed ###") # 放入None作为连接关闭的信号 ws.message_queue.put(None) def on_open(ws): def run(*args): for i in range(3): time.sleep(1) ws.send("Hello %d" % i) time.sleep(1) ws.close() print("thread terminating...") thread.start_new_thread(run, ()) def create_websocket_client(url, message_queue): websocket.enableTrace(True) ws = websocket.WebSocketApp(url, on_message=on_message, on_error=on_error, on_close=on_close) # 将队列绑定到WebSocket实例上 ws.message_queue = message_queue ws.on_open = on_open # 启动WebSocket客户端在独立线程 thread.start_new_thread(ws.run_forever, ()) return ws
主程序代码
import Queue import time # 假设上面的WebSocket代码存放在websocket_client.py文件中 from websocket_client import create_websocket_client if __name__ == "__main__": # 创建一个线程安全的消息队列 msg_queue = Queue.Queue() # 初始化WebSocket客户端,传入队列 ws = create_websocket_client("ws://<your-url>", msg_queue) # 主程序循环监听队列中的消息 while True: msg = msg_queue.get() if msg is None: # 收到None表示WebSocket连接关闭,退出循环 break # 这里添加你的主程序业务逻辑,比如处理消息、存储数据等 print("主程序收到WebSocket消息:", msg) # process_message(msg) # 替换为你的实际处理函数 print("主程序执行结束")
这个方案的优势在于消息传递可靠,线程安全,逻辑清晰,适合绝大多数场景。
2. 注入自定义消息处理回调
如果你希望主程序直接控制消息的处理逻辑,可以将主程序的消息处理函数注入到WebSocket客户端中,让on_message直接调用这个自定义函数。
WebSocket工具代码
try: import thread except ImportError: import _thread as thread import time import websocket def on_message(ws, message): # 调用主程序传入的自定义处理函数 ws.message_handler(message) def on_error(ws, error): print(error) def on_close(ws): print("### closed ###") # 可选:调用主程序传入的关闭回调 if hasattr(ws, 'close_handler'): ws.close_handler() def on_open(ws): def run(*args): for i in range(3): time.sleep(1) ws.send("Hello %d" % i) time.sleep(1) ws.close() print("thread terminating...") thread.start_new_thread(run, ()) def create_websocket_client(url, message_handler, close_handler=None): websocket.enableTrace(True) ws = websocket.WebSocketApp(url, on_message=on_message, on_error=on_error, on_close=on_close) # 绑定自定义处理函数 ws.message_handler = message_handler if close_handler: ws.close_handler = close_handler ws.on_open = on_open # 启动WebSocket客户端在独立线程 thread.start_new_thread(ws.run_forever, ()) return ws
主程序代码
import time from websocket_client import create_websocket_client def main_message_processor(message): # 主程序的消息处理逻辑,可根据需求自定义 print("主程序处理WebSocket消息:", message) # 示例:解析消息、更新状态、存储到数据库等 # parsed_data = parse_websocket_message(message) # update_application_state(parsed_data) def main_close_handler(): print("主程序收到WebSocket连接关闭通知") if __name__ == "__main__": # 初始化WebSocket客户端,传入自定义处理函数 ws = create_websocket_client("ws://<your-url>", message_handler=main_message_processor, close_handler=main_close_handler) # 主程序可以继续执行其他逻辑,比如保持运行或处理其他任务 try: while True: time.sleep(1) # 主程序的其他业务逻辑,比如定时任务、用户交互等 except KeyboardInterrupt: ws.close() print("主程序被用户中断")
这个方案的灵活性很高,主程序可以直接定义消息的处理流程,适合需要实时响应消息的场景。
3. 全局变量加线程锁(不推荐,仅作参考)
虽然可以通过全局变量共享消息,但多线程环境下必须使用线程锁来保证数据安全。这个方案逻辑繁琐,容易出现线程安全问题,仅作为备选参考。
WebSocket工具代码
try: import thread except ImportError: import _thread as thread import time import websocket import threading # Python 2.7的线程锁模块 # 全局变量存储消息 global_messages = [] # 线程锁,确保多线程访问全局变量的安全性 message_lock = threading.Lock() def on_message(ws, message): with message_lock: global_messages.append(message) def on_error(ws, error): print(error) def on_close(ws): print("### closed ###") def on_open(ws): def run(*args): for i in range(3): time.sleep(1) ws.send("Hello %d" % i) time.sleep(1) ws.close() print("thread terminating...") thread.start_new_thread(run, ()) def start_websocket_client(url): websocket.enableTrace(True) ws = websocket.WebSocketApp(url, on_message=on_message, on_error=on_error, on_close=on_close) ws.on_open = on_open thread.start_new_thread(ws.run_forever, ()) return ws
主程序代码
import time from websocket_client import global_messages, message_lock, start_websocket_client if __name__ == "__main__": ws = start_websocket_client("ws://<your-url>") while True: # 获取消息时必须加锁,避免数据竞争 with message_lock: if global_messages: # 取出所有消息并清空列表 received_messages = global_messages[:] global_messages.clear() else: received_messages = [] if received_messages: for msg in received_messages: print("主程序收到消息:", msg) # 处理消息的业务逻辑 time.sleep(0.5) # 可添加退出条件,比如检测WebSocket连接状态等
这个方案需要手动管理锁和全局变量,维护成本高,容易出现难以排查的线程问题,所以优先推荐前两个方案。
内容的提问来源于stack exchange,提问作者Abin
相关产品推荐
相关产品推荐

