Python跨进程共享下载任务链接的高效实现方案咨询
Python跨进程生产消费场景方案选型
你目前用fcntl.flock加锁共享文件的方案可以满足最基础的功能,但存在锁异常难排查、无任务状态追踪、扩容性差的问题,针对你这种生产者产出速率和消费者处理速率不匹配的场景,按推荐优先级给出以下更稳定、开发成本更低的方案:
方案1:基于Redis的任务队列(生产环境首选)
这是这类生产消费模型最通用的实现方案,天然适配速率差场景,不需要自己实现并发锁逻辑:
- 进程1(生产者)每15分钟查询数据库拿到待下载链接后,直接将链接推入Redis的List结构即可,队列本身会自动做数据缓冲,哪怕一次性推入数万条链接也不会丢失、不会阻塞。
- 进程2(消费者)用阻塞读指令拉取队列中的链接,拿到链接后执行下载逻辑即可,不需要轮询、不需要自己处理并发冲突。
- 后续如果下载压力增大,直接启动多个进程2实例即可,Redis天然支持多消费者竞争,不会出现重复消费的问题。
核心实现代码参考:
进程1(生产者)代码
import time import redis # 初始化Redis连接 r = redis.Redis(host="127.0.0.1", port=6379, db=0) QUEUE_KEY = "pending_artifact_download" def query_pending_links_from_db(): # 替换为你的数据库查询逻辑,返回待下载链接列表 return [ "https://your-internal-repo/artifact_v1.tar.gz", "https://your-internal-repo/artifact_v2.zip" ] if __name__ == "__main__": while True: pending_links = query_pending_links_from_db() for link in pending_links: r.lpush(QUEUE_KEY, link) # 间隔15分钟执行下一次查询 time.sleep(15 * 60)
进程2(消费者)代码
import os import redis import requests r = redis.Redis(host="127.0.0.1", port=6379, db=0) QUEUE_KEY = "pending_artifact_download" SAVE_PATH = "/data/artifact_storage" os.makedirs(SAVE_PATH, exist_ok=True) def download_artifact(url: str): file_name = url.rsplit("/", maxsplit=1)[-1] target_file = os.path.join(SAVE_PATH, file_name) with requests.get(url, stream=True, timeout=30) as resp: resp.raise_for_status() with open(target_file, "wb") as f: for chunk in resp.iter_content(chunk_size=1024*1024): f.write(chunk) if __name__ == "__main__": while True: # 阻塞等待队列新元素,0表示永久等待不超时 _, link_bytes = r.brpop(QUEUE_KEY, 0) link = link_bytes.decode("utf-8") try: download_artifact(link) except Exception as e: # 下载失败可根据需求推入重试队列、打日志 print(f"下载制品{link}失败,错误信息:{str(e)}")
注意:如果你的运行环境已经部署Redis,这个方案的开发量比自己实现文件加锁逻辑更小,不需要处理文件读写边界、锁粒度、异常恢复等边角问题。
方案2:基于SQLite的轻量任务表(零额外依赖首选)
如果你不想额外部署Redis这类服务,直接用Python标准库自带的SQLite做任务存储即可,稳定性远高于裸文件加锁:
- 新建一张任务表,字段包含自增ID、制品链接、任务状态(待下载/下载中/下载成功/下载失败)、创建时间。
- 进程1每次查询完数据库,直接向表中插入状态为「待下载」的链接记录即可。
- 进程2每次通过事务拉取1条状态为「待下载」的记录,立刻将状态更新为「下载中」并提交事务,再执行下载逻辑,下载完成后更新状态为成功即可。
SQLite本身的事务锁会自动处理并发访问控制,哪怕启动多个下载进程也不会抢到同一条任务,同时自带持久化能力,进程重启后未完成的任务不会丢失,还可以直接通过SQL语句查询任务进度、失败任务列表。
方案3:保留共享文件加锁的实现(仅适合临时小规模场景)
如果只是跑一次性临时任务、总任务量不超过100条,不想引入额外组件,你当前的方案也可以用,但需要注意两个核心坑:
- 锁粒度要控制到最小:写文件时加排他锁,写完立刻释放;读文件拿锁后,要立刻把已经读取走的链接从文件中移除,再释放锁,避免重复消费。
- 要做文件内容校验:如果进程在写文件中途崩溃,会导致文件格式损坏,下次启动时需要先校验文件内容合法性,避免读取出错。
内容的提问来源于stack exchange,提问作者Anushka
相关产品推荐
相关产品推荐

