如何在Python中让生成器函数与主循环并行运行
持续读取更新文件并并行处理匹配条目
问题背景
需要持续读取一个频繁更新的文件,将匹配特定正则表达式的条目提取后,逐个执行耗时操作。现有顺序实现(如下代码)存在明显缺陷:读取文件与处理条目串行执行,处理耗时操作时会暂停文件读取,无法及时追踪文件后续更新;且读取到文件末尾后会直接终止,无法持续监听新增内容。
现有顺序代码:
# 生成器函数:读取文件并返回匹配关键词的行 def read_file(keyword): with open("file.txt", "r") as file: for line in file: if keyword in line: # 返回匹配的行(去除首尾空白) yield line.strip() # 主循环:处理生成器返回的匹配行 keyword = "example" for matching_line in read_file(keyword): # 执行大量耗时操作 print(matching_line)
希望实现文件读取与条目处理的并行执行,已考虑线程方案,询问其他可行实现方式。
可行实现方案
1. 队列+多线程(轻量通用方案)
用queue.Queue作为读取线程与处理线程的缓冲:读取线程持续追踪文件更新,将匹配行放入队列;处理线程从队列取数据执行操作,两者完全并行互不阻塞。同时优化读取逻辑,实现类似tail -f的持续监听效果:
import time import queue import threading import re def file_reader(q, regex_pattern): pattern = re.compile(regex_pattern) with open("file.txt", "r") as f: # 先读取文件现有内容 for line in f: if pattern.search(line): q.put(line.strip()) # 持续监听文件新增内容 while True: new_line = f.readline() if new_line: if pattern.search(new_line): q.put(new_line.strip()) else: time.sleep(0.1) # 无新内容时短暂休眠,降低CPU占用 def line_processor(q): while True: line = q.get() if line is None: # 用None作为线程终止信号 break # 替换为实际的耗时操作 print(f"处理条目: {line}") q.task_done() if __name__ == "__main__": q = queue.Queue(maxsize=100) # 设置队列上限,避免内存溢出 regex = r"example" # 启动读取线程 reader_thread = threading.Thread(target=file_reader, args=(q, regex), daemon=True) reader_thread.start() # 可启动多个处理线程提升效率 processor_thread = threading.Thread(target=line_processor, args=(q,), daemon=True) processor_thread.start() # 保持主线程运行 try: while True: time.sleep(1) except KeyboardInterrupt: q.put(None) processor_thread.join() print("程序终止")
2. 异步IO方案(IO密集型操作首选)
如果处理操作是IO密集型(如网络请求、文件写入),用异步IO可避免线程切换开销,提升效率。需借助aiofiles实现异步文件读取:
import asyncio import re from aiofiles import open async def file_reader(q, regex_pattern): pattern = re.compile(regex_pattern) async with open("file.txt", "r") as f: # 读取现有内容 async for line in f: if pattern.search(line): await q.put(line.strip()) # 持续监听新增内容 while True: line = await f.readline() if line: if pattern.search(line): await q.put(line.strip()) else: await asyncio.sleep(0.1) async def line_processor(q): while True: line = await q.get() if line is None: break # 替换为异步耗时操作 print(f"处理条目: {line}") q.task_done() async def main(): q = asyncio.Queue(maxsize=100) regex = r"example" # 创建异步任务 reader_task = asyncio.create_task(file_reader(q, regex)) processor_task = asyncio.create_task(line_processor(q)) try: await asyncio.gather(reader_task, processor_task) except KeyboardInterrupt: await q.put(None) await processor_task print("程序终止") if __name__ == "__main__": asyncio.run(main())
3. 多进程方案(CPU密集型操作首选)
如果处理操作是CPU密集型(如复杂计算、数据加密),多进程可绕过Python GIL限制,充分利用多核CPU资源:
import time import re from multiprocessing import Process, Queue def file_reader(q, regex_pattern): pattern = re.compile(regex_pattern) with open("file.txt", "r") as f: for line in f: if pattern.search(line): q.put(line.strip()) while True: new_line = f.readline() if new_line: if pattern.search(new_line): q.put(new_line.strip()) else: time.sleep(0.1) def line_processor(q): while True: line = q.get() if line is None: break # 替换为实际的CPU密集型操作 result = line.upper() * 1000 print(f"处理结果长度: {len(result)}") if __name__ == "__main__": q = Queue(maxsize=100) regex = r"example" reader_process = Process(target=file_reader, args=(q, regex), daemon=True) reader_process.start() processor_process = Process(target=line_processor, args=(q,), daemon=True) processor_process.start() try: while True: time.sleep(1) except KeyboardInterrupt: q.put(None) processor_process.join() print("程序终止")
方案选择建议
- IO密集型操作:优先选异步IO或线程方案,线程实现更简单直观
- CPU密集型操作:必须选多进程方案,避免GIL对性能的限制
- 所有方案都建议设置队列上限,防止文件更新过快导致内存溢出
内容的提问来源于stack exchange,提问作者Inc0gnito
相关产品推荐
相关产品推荐

