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

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)

核心原理

  1. PushAdapter内部维护一个预激的协程_receiver,初始状态停在第一个yield处等待接收数据。
  2. 调用send(item)时,协程接收数据并走到第二个yield,将数据暴露给迭代器。
  3. 迭代器的__next__方法通过next(self._coro)获取推送的数据,之后再调用一次next让协程回到接收等待状态,准备下一次推送。

内容的提问来源于stack exchange,提问作者Nat Goodspeed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 16:00:12