【问题标题】:Irregular Transmission Problem with Python Twisted Push ProducerPython Twisted Push Producer 的不规则传输问题
【发布时间】:2013-06-29 00:19:01
【问题描述】:

我想使用 Twisted 从队列传输数据。我目前使用推送生产者来轮询队列中的项目并写入传输。

class Producer:

    implements(interfaces.IPushProducer)

    def __init__(self, protocol, queue):
        self.queue = queue
        self.protocol = protocol

    def resumeProducing(self):
        self.paused = False
        while not self.paused:
            try:
                data = self.queue.get_nowait()
                logger.debug("Transmitting: '%s'", repr(data))
                data = cPickle.dumps(data)
                self.protocol.transport.write(data + "\r\n")
            except Empty:
                pass

    def pauseProducing(self):
        logger.debug("Transmitter paused.")
        self.paused = True

    def stopProducing(self):
        pass

问题是,数据发送非常不规则,如果队列中只有一项,则永远不会发送数据。似乎 Twisted 一直等到要传输的数据增长到特定值,直到它传输它。 我实施生产者的方式是否正确?我可以现在强制 Twisted 传输数据吗?

我也尝试过使用 pull producer,但 Twisted 根本不调用它的 resumeProducing() 方法。使用拉取生产者时,是否必须从外部调用 resumeProducer() 方法?

【问题讨论】:

    标签: python twisted producer-consumer


    【解决方案1】:

    如果没有看到完整的示例(也就是说,没有看到向消费者注册它的代码以及将项目放入该队列的代码),很难说为什么您的生产者不能很好地工作。

    但是,您可能会遇到的一个问题是,如果在调用resumeProducing 时您的队列是,那么您根本不会向消费者写入任何字节。当项目被放入队列时,它们将永远坐在那里,因为消费者不会再次调用您的 resumeProducing 方法。

    这可以概括为队列中没有足够数据导致消费者在生产者上调用pauseProducing 的任何其他情况。作为推送生产者,您的工作是继续自己生产数据,直到消费者致电pauseProducing(或stopProducing)。

    对于这种特殊情况,这可能意味着无论何时您要将某些东西放入该队列 - 停止:检查生产者是否没有暂停,如果没有,则将其写入消费者 而是。仅在生产者暂停时才将项目放入队列中。

    【讨论】:

    • 如果我必须从外部检查它是否正在运行,然后我才能给他处理数据,这不是对生产者干净设计的破坏吗?在我看来,数据本身的处理应该对外部透明。
    • 这是一个非常大的问题。一个答案是 - 是的,当然,它是;因此,您应该向生产者添加一个方法来将数据排队,并且该方法应该检查状态以确定它是否应该真正进入队列,或者是否应该立即将其写入消费者。这将所有逻辑保留在生产者类中。但是,考虑到如果这个生产者的所有用户都决定不知道它是否被暂停,那么他们最终可能会溢出队列。生产者的要点是传播低级缓冲区满事件,以避免溢出。
    • 好的,谢谢您的回答。似乎这是唯一的解决方案,因为没有发布其他答案。但是,我仍然对这个解决方案不满意。
    • 这个想法是你真的不应该“从外面检查”;将项目放入队列的东西也应该能够暂停。您应该将此通知传递到生成数据的责任链尽可能远的位置。在您无法进一步转发它的任何地方,您可能会被要求缓冲无限量的数据;此时,您需要考虑某种策略来处理队列中条目的无限增长:将它们保存到磁盘?向您的同伴发送拒绝消息?这取决于应用程序。
    【解决方案2】:

    这里有两种可能的解决方案:

    1) 定期轮询您的本地应用程序以查看您是否有其他数据要发送。

    注意。这依赖于twisted 中deferLater 方法的周期性异步回调。如果您需要按需发送数据的响应式应用程序,或者需要长时间运行的阻塞操作(例如,使用自己的事件循环的 ui),则可能不合适。

    代码:

    from twisted.internet.protocol import Factory
    from twisted.internet.endpoints import TCP4ServerEndpoint
    from twisted.internet.interfaces import IPushProducer
    from twisted.internet.task import deferLater, cooperate
    from twisted.internet.protocol import Protocol
    from twisted.internet import reactor
    from zope.interface import implementer
    import time
    
    # Deferred action
    def periodically_poll_for_push_actions_async(reactor, protocol):
      while True:
        protocol.send(b"Hello World\n")
        yield deferLater(reactor, 2, lambda: None)
    
    # Push protocol
    @implementer(IPushProducer)
    class PushProtocol(Protocol):
    
       def connectionMade(self):
         self.transport.registerProducer(self, True)
         gen = periodically_poll_for_push_actions_async(self.transport.reactor, self)
         self.task = cooperate(gen)
    
       def dataReceived(self, data):
         self.transport.write(data)
    
       def send(self, data):
         self.transport.write(data)
    
       def pauseProducing(self):
         print 'Workload paused'
         self.task.pause()
    
       def resumeProducing(self):
         print 'Workload resumed'
         self.task.resume()
    
       def stopProducing(self):
         print 'Workload stopped'
         self.task.stop()
    
       def connectionLost(self, reason):
         print 'Connection lost'
         try:
           self.task.stop()
         except:
           pass
    
    # Push factory
    class PushFactory(Factory):
      def buildProtocol(self, addr):
        return PushProtocol()
    
    # Run the reactor that serves everything
    endpoint = TCP4ServerEndpoint(reactor, 8089)
    endpoint.listen(PushFactory())
    reactor.run()
    

    2) 手动跟踪协议实例并使用来自不同线程的 reactor.callFromThread()。让您摆脱另一个线程中的长时间阻塞操作(例如 ui 事件循环)。

    代码:

    from twisted.internet.protocol import Factory
    from twisted.internet.endpoints import TCP4ServerEndpoint
    from twisted.internet.interfaces import IPushProducer
    from twisted.internet.task import deferLater, cooperate
    from twisted.internet.protocol import Protocol
    from twisted.internet import reactor, threads
    import time
    import random
    import threading
    
    # Connection
    protocol = None
    
    # Some other thread that does whatever it likes.
    class SomeThread(threading.Thread):
      def run(self):
        while True:
          print("Thread loop")
          time.sleep(random.randint(0, 4))
          if protocol is not None:
            reactor.callFromThread(self.dispatch)
      def dispatch(self):
        global protocol
        protocol.send("Hello World\n")
    
    # Push protocol
    class PushProtocol(Protocol):
    
       def connectionMade(self):
         global protocol
         protocol = self
    
       def dataReceived(self, data):
         self.transport.write(data)
    
       def send(self, data):
         self.transport.write(data)
    
       def connectionLost(self, reason):
         print 'Connection lost'
    
    # Push factory
    class PushFactory(Factory):
      def buildProtocol(self, addr):
        return PushProtocol()
    
    # Start thread
    other = SomeThread()
    other.start()
    
    # Run the reactor that serves everything
    endpoint = TCP4ServerEndpoint(reactor, 8089)
    endpoint.listen(PushFactory())
    reactor.run()
    

    就我个人而言,我发现 IPushProducer 和 IPullProducer 需要定期回调,这使得它们的用处不大。其他人不同意... 耸耸肩。任君挑选。

    【讨论】:

    • 您似乎不了解流控制和缓冲区管理的工作原理。如果没有来自传输的 IConsumer 实现的通知(带有 registerProducer 方法的东西,调用 pauseProducing/resumeProducing 的东西),如果您的一个客户阅读速度太慢(这被称为“涓流攻击”,如果是故意的,或者是“糟糕的 GSM 连接”),那么传递给 transport.write 的数据将无限缓冲,您的程序最终将耗尽内存并崩溃。
    • @Glyph 如果您对如何编写推送消费者有更好的解决方案,请将其作为解决方案发布。这里有两个事实:1)你不能编写一个不经常使用 IPushProducer 推送更新的服务器。 2) 你不能用 IPullProducer 写一个。所以,如果你有一些聪明的方法,发布它怎么样?
    • @Jean-PaulCalderone 不,你不能。如果您发送无数据的时间不确定,IPushProducer 将永远不会被挂起。不这么认为?继续,发布一个解决方案。请记住,除非调用 transport.write(),否则您将永远不会被挂起。
    • 如果你不发送数据,你不需要被“暂停”(我们通常说暂停)。您只需在下游发送缓冲区已满时暂停。这是一个例子 - gist.github.com/exarkun/5884100
    • 我的立场是正确的。 ...尽管,您的要点似乎不适用于 twisted-13.0.0。不过,我已经更新了我的答案,以反映延迟 IO 是(可能)可行的替代方案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-05-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-16
    相关资源
    最近更新 更多