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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 19:05:28