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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 10:06:20