itertools.groupby分批次处理后如何反向合并为单元素生成器?
问题原因
你当前拿到的元素是生成器/迭代器的核心原因是:executor.map 本身返回的是可迭代对象,你原有写法的chain(batches)仅展开了最外层的批次容器,没有进一步展开每个批次内部的迭代器元素,因此遍历得到的i仍然是每个批次对应的迭代器对象。
另外你当前的thread函数存在隐藏风险:ThreadPoolExecutor的with上下文管理器会在代码块退出时自动关闭线程池,而executor.map是惰性求值的,如果你在上下文外再迭代其返回的迭代器,会抛出线程池已关闭的错误。
解决方案
1. 修复flat_map实现
两种实现方式任选即可:
方式一:使用itertools.chain.from_iterable(推荐)
专门用于展开嵌套一层的可迭代对象,代码更简洁:
from itertools import chain from typing import Iterable def flat_map(batches: Iterable): yield from chain.from_iterable(batches)
方式二:手动嵌套遍历(更易理解)
from typing import Iterable def flat_map(batches: Iterable): for batch in batches: for item in batch: yield item
2. 修复thread函数的潜在问题
在yield前将executor.map的结果转为列表,保证任务在线程池关闭前执行完成:
from concurrent.futures import ThreadPoolExecutor from typing import Iterable, Callable def thread(data: Iterable, func: Callable, n=4): with ThreadPoolExecutor(max_workers=n) as executor: for batch in data: # 转为列表触发实际执行,避免后续迭代时线程池已关闭 yield list(executor.map(func, batch))
完整测试示例
from typing import List, Any, Iterable, Callable from itertools import groupby, count, chain from concurrent.futures import ThreadPoolExecutor def batch(data: List[Any], size=4): c = count() for _, g in groupby(data, lambda _: next(c)//size): yield g def thread(data: Iterable, func: Callable, n=4): with ThreadPoolExecutor(max_workers=n) as executor: for batch in data: yield list(executor.map(func, batch)) def flat_map(batches: Iterable): yield from chain.from_iterable(batches) # 测试用函数 def square(x: int) -> int: return x * x if __name__ == "__main__": raw_data = list(range(10)) # 完整数据流管道 processed_data = flat_map(thread(batch(raw_data, size=3), square, n=2)) # 转为列表触发执行 print(list(processed_data)) # 输出:[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
内容的提问来源于stack exchange,提问作者CpILL
相关产品推荐
相关产品推荐

