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

如何实现从API请求到SQLite数据库的多进程处理?

多进程并行API请求+串行SQLite上传方案

核心思路

用生产者-消费者模式拆分任务,兼顾效率与数据安全:

  • 生产者:多进程并行发起API请求,获取并转换数据后放入进程安全队列
  • 消费者:单进程从队列取数据,串行写入SQLite,彻底避免多进程写库冲突,保证数据不丢失

适用模块

  • multiprocessing:提供进程创建、进程安全队列(Queue),实现跨进程数据传递
  • requests:发起API请求
  • pandas:处理JSON格式的API返回数据
  • sqlite3:数据库操作(注意:每个进程必须独立创建连接,绝对不能跨进程共享连接)

代码实现

import sqlite3
import requests
import pandas as pd
from multiprocessing import Process, Queue, cpu_count

# 替换为你的实际配置
URL_LIST = ["url1", "url2", "url3"]
DB_PATH = "testDB"
END_SIGNAL = None  # 用于通知消费者所有数据已处理完毕

def producer(url, queue):
    """生产者进程:请求API并把处理后的数据放入队列"""
    try:
        r = requests.get(url)
        r.raise_for_status()  # 捕获HTTP请求异常
        data = pd.read_json(r.text)
        # 转成字典列表,方便消费者批量处理
        queue.put(data.to_dict("records"))
    except Exception as e:
        print(f"请求URL {url} 失败: {str(e)}")

def consumer(queue):
    """消费者进程:从队列取数据,串行写入SQLite"""
    # 消费者必须自己创建数据库连接,不能用其他进程的连接
    con = sqlite3.connect(DB_PATH)
    cur = con.cursor()
    # 提前创建表(如果不存在)
    cur.execute("CREATE TABLE IF NOT EXISTS test (valA INTEGER PRIMARY KEY, valB TEXT);")
    con.commit()

    while True:
        item = queue.get()
        if item == END_SIGNAL:
            break  # 收到结束信号,退出循环
        if not item:
            continue
        
        # 批量插入提升效率,同时保证事务原子性
        try:
            cur.executemany(
                "INSERT OR REPLACE INTO test (valA, valB) VALUES (?, ?);",
                [(row["valA"], row["valB"]) for row in item]
            )
            con.commit()
            print(f"成功上传 {len(item)} 条数据")
        except Exception as e:
            con.rollback()  # 出错回滚,避免部分数据写入
            print(f"数据上传失败: {str(e)}")
    con.close()

if __name__ == "__main__":
    # 创建进程安全队列
    data_queue = Queue()

    # 启动消费者进程
    consumer_proc = Process(target=consumer, args=(data_queue,))
    consumer_proc.start()

    # 启动生产者进程(控制并发数,避免触发API限流)
    producer_procs = []
    max_concurrent_procs = min(cpu_count(), len(URL_LIST))  # 按CPU核心数或URL数量取小值
    for url in URL_LIST:
        proc = Process(target=producer, args=(url, data_queue))
        producer_procs.append(proc)
        proc.start()
        
        # 达到最大并发数时,等待已有进程完成再启动新的
        if len(producer_procs) >= max_concurrent_procs:
            for p in producer_procs:
                p.join()
            producer_procs = []

    # 等待剩余生产者进程结束
    for p in producer_procs:
        p.join()

    # 给消费者发结束信号
    data_queue.put(END_SIGNAL)
    # 等待消费者处理完所有数据
    consumer_proc.join()

    print("所有数据处理完成")

关键注意事项

  • 数据库连接隔离:SQLite的连接对象不能跨进程共享,每个进程必须自己创建连接,否则会导致数据损坏或程序崩溃。
  • 进程安全队列:必须用multiprocessing.Queue,普通的queue.Queue只适用于线程通信,进程间通信必须用专门的进程安全队列。
  • 原子性保障:用executemany批量插入+事务回滚,确保一批数据要么全部写入成功,要么全部回滚,不会出现部分数据丢失的情况。
  • 并发控制:生产者进程数不要设置过高,避免触发API的频率限制,建议参考CPU核心数或API官方给出的并发限制。
  • 异常捕获:给请求和上传环节加异常处理,避免单个任务失败导致整个程序挂掉。

内容的提问来源于stack exchange,提问作者flowerboy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 23:30:49