Python中如何向被其他函数消费的迭代器逐个推送元素?
问题描述
现有一个接收可迭代对象的函数,比如:
def printall(iterable): for item in iterable: print(item)
要求不修改、不克隆该函数,希望通过异步推送的方式给它提供数据,理想的伪代码如下:
push_handle = XXXX() printall(push_handle) # printall() 此时等待 next(iter(iterable)) for item0 in some_source(): # 一些复杂处理逻辑 statement statement item1 = some_transformation(item0) push_handle.send(item1) # 发送后printall()完成一次迭代,继续等待下一个数据
虽然可以把处理逻辑重构为生成器来实现:
def processing(iterable): for item0 in iterable: statement statement item1 = some_transformation(item0) yield item1 printall(processing(some_source()))
但希望能实现伪代码中的“推送式适配器”,避免重构代码。尝试过用生成器,但生成器的yield是返回给直接调用者,无法实现这种跨控制流的推送。考虑过线程+队列,但太重;也试过eventlet的无界队列,但需要分属不同协程。感觉应该有更轻量的方式,是不是忽略了什么?
解决方案
可以用原生协程+迭代器结合的方式实现这个推送适配器,完全基于Python内置机制,无需额外依赖或线程/协程框架(如果消费函数是阻塞式迭代,仅需简单线程分离生产/消费逻辑)。
适配器实现
class PushAdapter: def __init__(self): # 初始化接收协程,预激后停在第一个yield等待数据 self._coro = self._receiver() next(self._coro) def _receiver(self): while True: # 接收send过来的数据,挂起等待迭代器取走 item = yield # 将数据传递给迭代器的__next__方法 yield item def send(self, item): # 向协程推送数据 return self._coro.send(item) def __iter__(self): return self def __next__(self): # 迭代器获取数据时,驱动协程拿到推送内容 result = next(self._coro) # 让协程回到接收数据的等待状态 next(self._coro) return result
使用示例
def printall(iterable): for item in iterable: print(f"printall 输出: {item}") def some_source(): # 模拟原始数据源 yield "raw_data_1" yield "raw_data_2" yield "raw_data_3" def some_transformation(item): return f"processed_{item}" # 初始化适配器 push_handle = PushAdapter() # 启动printall(用守护线程分离消费逻辑,避免阻塞生产流程) import threading threading.Thread(target=printall, args=(push_handle,), daemon=True).start() # 模拟生产/推送流程 for item0 in some_source(): print(f"正在处理原始数据: {item0}") # 模拟复杂处理步骤 item1 = some_transformation(item0) push_handle.send(item1)
核心原理
PushAdapter内部维护一个预激的协程_receiver,初始状态停在第一个yield处等待接收数据。- 调用
send(item)时,协程接收数据并走到第二个yield,将数据暴露给迭代器。 - 迭代器的
__next__方法通过next(self._coro)获取推送的数据,之后再调用一次next让协程回到接收等待状态,准备下一次推送。
内容的提问来源于stack exchange,提问作者Nat Goodspeed
相关产品推荐
相关产品推荐

