【问题标题】:Fire DeferredList after multiple callbacks多次回调后触发 DeferredList
【发布时间】:2018-04-25 08:05:51
【问题描述】:

twisted.internet.defer.DeferredList 这样做:

我将一组延迟合并到一个回调中。

我为他们的回调跟踪了一个 Deferred 列表,并制作了一个 全部完成后回调,(成功,结果)列表 元组,“成功”是一个布尔值。

请注意,在将 Deferred 放入 延迟列表。例如,您可以在 通过将 errbacks 添加到 Deferreds after put 来延迟'消息 它们在 DeferredList 中,因为 DeferredList 不会吞下错误。 (虽然更方便的方法是简单地设置 消费错误标志)

def __init__(self, deferredList, fireOnOneCallback=0, fireOnOneErrback=0, consumeErrors=0): (source)
    overrides twisted.internet.defer.Deferred.__init__
    Initialize a DeferredList.
    Parameters  deferredList    The list of deferreds to track. (type: list of Deferreds )
    fireOnOneCallback   (keyword param) a flag indicating that only one callback needs to be fired for me to call my callback
    fireOnOneErrback    (keyword param) a flag indicating that only one errback needs to be fired for me to call my errback
    consumeErrors   (keyword param) a flag indicating that any errors raised in the original deferreds should be consumed by this DeferredList. This is useful to prevent spurious warnings being logged.

具体来说:

fireOnOneCallback(关键字参数)一个标志,表明只有一个 需要触发回调才能调用我的回调

我正在寻找像fireOnOneCallback=True 这样的行为,而是在n 回调上触发。我试过这样做,但它已经变成了一团糟。我相信有更好的方法。

def _get_fired_index(deferred_list):
    for index, (success, value) in enumerate(deferred_list):
        if success:
            return index
    raise ValueError('No deferreds were fired.')


def _fire_on_other_callback(already_fired_index, deferred_list, callback, ):
    dlist_except_first_fired = (
        deferred_list[:already_fired_index]
        + deferred_list[already_fired_index + 1:]
    )
    dlist2 = DeferredList(dlist_except_first_fired, fireOnOneCallback=True)
    dlist2.addCallback(callback, deferred_list)


def _fire_on_two_callbacks(deferreds, callback, errback):
    dlist1 = DeferredList(deferreds, fireOnOneCallback=True)
    dlist1.addCallback(_get_fired_index)
    dlist1.addCallback(_fire_on_other_callback, deferreds, callback, errback)

【问题讨论】:

    标签: python twisted


    【解决方案1】:

    这是一种可能的方法。

    from __future__ import print_function
    
    import attr
    from twisted.internet.defer import Deferred
    
    def fireOnN(n, ds):
        acc = _Accumulator(n)
        for index, d in enumerate(ds):
            d.addCallback(acc.one_result, index)
        return acc.n_results
    
    @attr.s
    class _Accumulator(object):
        n = attr.ib()
        so_far = attr.ib(default=attr.Factory(dict))
        done = attr.ib(default=False)
        n_results = attr.ib(default=attr.Factory(Deferred))
    
        def one_result(self, result, index):
            if self.done:
                return result
            self.so_far[index] = result
            if len(self.so_far) == self.n:
                self.done = True
                so_far = self.so_far
                self.so_far = None
                self.n_results.callback(so_far)
    
    dx = list(Deferred().addCallback(print, i) for i in range(3))
    done = fireOnN(2, dx)
    done.addCallback(print, "done")
    
    for i, d in enumerate(dx):
        d.callback("result {}".format(i))
    

    请注意,此实现不处理 errbacks 并且可能还有其他缺点(例如坚持使用 n_results 参考)。但是,基本思想是合理的:从回调中累积状态,直到达到所需的条件,然后触发另一个 Deferred。

    DeferredList 只会给这个问题带来不必要的复杂性,其无关的功能和界面并不是为解决这个问题而设计的。

    【讨论】:

      【解决方案2】:

      这是使用DeferredSemaphore 处理潜在竞争条件的另一种方法。这将在n 的延迟对象被触发并取消其余部分后立即触发。

      from twisted.internet import defer
      
      
      def fireAfterNthCallback(deferreds, n):
          if not n or n > len(deferreds):
              raise ValueError
      
          results = {}
          finished_deferred = defer.Deferred()
          sem = defer.DeferredSemaphore(1)
      
          def wrap_sem(result, index):
              return sem.run(callback_result, result, index)
      
          def cancel_remaining():
              finished = [deferreds[index] for index in results.keys()]
              for d in finished:
                  deferreds.remove(d)
              for d in deferreds:
                  d.addErrback(lambda err: err.trap(defer.CancelledError))
                  d.cancel()
      
          def callback_result(result, index):
              results[index] = result
              if len(results) >= n:
                  cancel_remaining()
                  finished_deferred.callback(results.values())
              return result
      
          for deferred_index, deferred in enumerate(deferreds):
              deferred.addCallback(wrap_sem, deferred_index)
      
          return finished_deferred
      

      【讨论】:

        猜你喜欢
        • 2015-01-30
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多