基于Twisted的API发送线程启动后仅执行一次循环问题求助
问题
我正尝试基于Twisted开发API(依赖项目:OpenApiPy)。以下是一段启动twisted reactor的不完整代码:
def onError(failure): print("Message Error: ", failure) def connected(client): print("\nConnected") request = ProtoOAApplicationAuthReq() request.clientId = app_id request.clientSecret = app_secret deferred = client.send(request) deferred.addErrback(onError) def disconnected(client, reason): # Callback for client disconnection print("\nDisconnected: ", reason) def onMessageReceived(client, message): # Callback for receiving all messages print("Message received: \n", Protobuf.extract(message)) def send_messages(): print('waiting for messages to send') while True: print(1) message = message_queue.get() deferred = client.send(message) deferred.addErrback(onError) time.sleep(10) def send_message_thread(message): message_queue.put(message) def main(): client.setConnectedCallback(connected) client.setDisconnectedCallback(disconnected) client.setMessageReceivedCallback(onMessageReceived) client.startService() # send_messages() send_thread = threading.Thread(target=send_messages) send_thread.daemon = True # The thread will exit when the main program exits send_thread.start() # Run Twisted reactor reactor.run() if __name__ == "__main__": main()
程序成功连接后输出如下:
waiting for messages to send 1 Connected Message received: Message received: Message received: Message received: ...
可见send_thread已启动,但其中的while True循环仅执行一次(仅打印一次1)。我猜测可能被reactor.run()阻塞,但对Twisted不熟悉,请问该问题原因是什么?
原因及解决办法
问题出在message_queue.get()这一行——标准库queue.Queue.get()默认是阻塞式调用,当队列中没有待发送的消息时,它会直接挂起当前线程,直到有新消息被放入队列才会继续执行。你的输出里打印了一次1,说明循环执行到message_queue.get()就停住了,和reactor.run()没有关系。
另外还有两个不符合Twisted设计的潜在问题:
- 在外部线程中调用
client.send():Twisted是单线程事件驱动框架,绝大多数API都不是线程安全的,直接在非Reactor线程调用可能引发竞态条件或未知异常。 - 使用
time.sleep(10):这是同步阻塞操作,会完全占用子线程,同时违背了Twisted的协作式调度逻辑。
正确的做法是完全基于Twisted的异步机制重构代码:
- 用Twisted提供的
twisted.python.queue.Queue替代标准库队列,支持异步获取消息 - 用
reactor.callLater实现定时逻辑,替代while True+sleep - 所有操作都在Reactor线程中执行,避免线程安全问题
示例修改后的核心代码:
from twisted.python.queue import Queue from twisted.internet import reactor # 替换为Twisted异步队列 message_queue = Queue() def send_messages(): def _handle_next_message(): # 异步获取队列中的消息,不会阻塞Reactor def _send(msg): deferred = client.send(msg) deferred.addErrback(onError) # 10秒后继续处理下一条消息 reactor.callLater(10, _handle_next_message) message_queue.get().addCallback(_send).addErrback(lambda err: print(err)) # 启动消息循环 _handle_next_message() def send_message(message): # 异步放入队列 message_queue.put(message) def main(): client.setConnectedCallback(connected) client.setDisconnectedCallback(disconnected) client.setMessageReceivedCallback(onMessageReceived) client.startService() # 启动异步消息发送逻辑,无需额外线程 send_messages() reactor.run()
这样既符合Twisted的异步编程模型,也彻底解决了线程阻塞和安全问题。
内容的提问来源于stack exchange,提问作者wiktor
相关产品推荐
相关产品推荐

