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

基于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设计的潜在问题:

  1. 在外部线程中调用client.send():Twisted是单线程事件驱动框架,绝大多数API都不是线程安全的,直接在非Reactor线程调用可能引发竞态条件或未知异常。
  2. 使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 01:20:54