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

Python多进程批量Steam API调用提速及数据库插入优化问询

看起来你现在面临的是Steam用户数据批量抓取的效率瓶颈,7-8亿条数据按当前速度要跑几百天确实太夸张了。我帮你分析下现有代码的问题,然后给出能大幅提速的优化方案:

优化Steam用户数据批量抓取与入库方案

先说说你现有代码里拖慢速度的核心问题

  • API请求极度浪费:每次只查一个SteamID,但Steam的GetPlayerSummaries接口一次支持查最多100个!这直接能把请求次数砍到原来的1%,效率提升巨大。
  • 线程/进程混用逻辑错误:你在databaseFiller里递归创建线程池并调用自身,这会导致无限递归,直接耗尽系统资源,完全是逻辑bug。
  • 数据库插入效率极低:逐条插入数据,哪怕用了队列还是循环单条执行,数据库IO本身就慢,这么搞肯定拖后腿。
  • 内存去重完全不可行:全局的image_hashes在多进程环境下是各进程独立的副本,而且7亿条哈希值根本不可能存在内存里,应该让数据库用唯一约束自动去重。
  • 数据库连接管理混乱:多进程共享一个数据库连接会引发各种冲突,每个进程应该使用独立的连接。

针对性的优化方案,一步一步来

1. 批量发起API请求

Steam的GetPlayerSummaries接口支持用逗号分隔一次传入最多100个steamids,把你的SteamID后缀分成每100个一组批量请求,这是效率提升的核心。

2. 合理搭配进程池与线程池

  • 用进程池处理不同的SteamID区间(进程能利用多核CPU资源),每个进程内部用线程池处理API请求和图片下载(这类IO密集型任务适合用线程池)。
  • 彻底删掉递归调用的逻辑:每个进程负责一个区间,把区间拆成小批量,批量请求API、处理数据、批量插入数据库,流程清晰不混乱。

3. 数据库层面优化

  • 给hash字段添加唯一约束,让数据库自动跳过重复数据,不用在内存里存所有哈希值:
    CREATE TABLE IF NOT EXISTS user (
        id INTEGER PRIMARY KEY AUTOINCREMENT,
        hash TEXT UNIQUE,
        profile TEXT
    );
    
  • 使用executemany批量插入数据,一次插几百条甚至几千条,比逐条插入效率提升几十倍。
  • 每个进程创建独立的数据库连接,避免多进程共享连接的冲突问题。

4. 添加速率限制与重试机制

Steam API有请求频率限制,别猛刷导致被封禁。每批请求之间加个小间隔(比如1秒),同时对失败的请求进行重试,避免数据丢失。

5. 图片处理优化

  • 先判断是否是默认头像,是的话直接跳过,不用下载无用图片。
  • 图片下载与哈希计算可以放到线程池里,和API请求并行处理,节省整体耗时。

修改后的示例代码

import traceback
import requests
import sys
from time import sleep
from multiprocessing import Pool
from io import BytesIO
import imagehash
from PIL import Image
import sqlite3
from itertools import islice

# 配置参数
MIN_STEAMID_SUFFIX = 7960265729
MAX_STEAMID_SUFFIX = 9080098567
DATABASE_LOCATION = 'D:/Script/steam_database.db'
API_KEY = "你的Steam API密钥"  # 替换成实际密钥
PROCESS_POOL_SIZE = 8  # 根据CPU核心数调整,别设太大
BATCH_SIZE_API = 100  # 每次API请求的SteamID数量
BATCH_SIZE_DB = 1000  # 每次数据库批量插入的数量
RATE_LIMIT = 1  # 每批请求间隔(秒),按需调整

def init_db():
    # 初始化数据库,创建带唯一约束的表
    conn = sqlite3.connect(DATABASE_LOCATION)
    cursor = conn.cursor()
    cursor.execute("""
        CREATE TABLE IF NOT EXISTS user (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            hash TEXT UNIQUE,
            profile TEXT
        )
    """)
    conn.commit()
    conn.close()

def chunk_iterable(iterable, size):
    # 将可迭代对象拆分成指定大小的块
    iterator = iter(iterable)
    while chunk := list(islice(iterator, size)):
        yield chunk

def process_steamid_batch(steamid_suffixes):
    # 处理一批SteamID后缀,请求API并处理数据
    steamids = [f"7656119{suffix}" for suffix in steamid_suffixes]
    steamids_str = ",".join(steamids)
    
    # 请求API,带重试逻辑
    retry_count = 3
    for _ in range(retry_count):
        try:
            response = requests.get(
                f'http://api.steampowered.com/ISteamUser/GetPlayerSummaries/v0002/?key={API_KEY}&steamids={steamids_str}'
            )
            response.raise_for_status()
            data = response.json()
            break
        except Exception as e:
            print(f"API请求失败,重试中: {e}")
            sleep(2)
    else:
        print(f"API请求重试{retry_count}次失败,跳过该批次")
        return []
    
    players = data.get('response', {}).get('players', [])
    processed_data = []
    default_avatar = "https://steamcdn-a.akamaihd.net/steamcommunity/public/images/avatars/fe/fef49e7fa7e1997310d705b2a6158ff8dc1cdfeb.jpg"
    
    for player in players:
        pfp_url = player.get('avatar')
        profile_url = player.get('profileurl')
        
        if not pfp_url or pfp_url == default_avatar or not profile_url:
            continue
        
        # 下载图片并计算哈希
        try:
            img_response = requests.get(pfp_url)
            img_response.raise_for_status()
            img = Image.open(BytesIO(img_response.content))
            img_hash = str(imagehash.average_hash(img))
            processed_data.append((img_hash, profile_url))
        except Exception as e:
            print(f"处理图片失败: {e}")
            continue
    
    return processed_data

def worker_process(suffix_range):
    # 单个工作进程处理一个SteamID后缀区间
    start_suffix, end_suffix = suffix_range
    print(f"开始处理区间: {start_suffix} → {end_suffix}")
    
    # 每个进程创建独立的数据库连接
    conn = sqlite3.connect(DATABASE_LOCATION)
    cursor = conn.cursor()
    
    # 拆分区间为API批次处理
    suffix_iter = range(start_suffix, end_suffix)
    for batch_suffixes in chunk_iterable(suffix_iter, BATCH_SIZE_API):
        batch_data = process_steamid_batch(batch_suffixes)
        
        if not batch_data:
            sleep(RATE_LIMIT)
            continue
        
        # 按数据库批量大小拆分插入
        for db_batch in chunk_iterable(batch_data, BATCH_SIZE_DB):
            try:
                cursor.executemany(
                    "INSERT OR IGNORE INTO user (hash, profile) VALUES (?, ?);",
                    db_batch
                )
                conn.commit()
                print(f"成功插入 {len(db_batch)} 条数据")
            except Exception as e:
                print(f"数据库插入失败: {e}")
                conn.rollback()
        
        sleep(RATE_LIMIT)  # 速率限制
    
    conn.close()
    print(f"完成处理区间: {start_suffix} → {end_suffix}")

if __name__ == '__main__':
    init_db()
    
    # 拆分SteamID后缀区间到各个进程
    total_suffixes = MAX_STEAMID_SUFFIX - MIN_STEAMID_SUFFIX
    suffixes_per_process = total_suffixes // PROCESS_POOL_SIZE
    ranges = []
    
    for i in range(PROCESS_POOL_SIZE):
        start = MIN_STEAMID_SUFFIX + i * suffixes_per_process
        end = MAX_STEAMID_SUFFIX if i == PROCESS_POOL_SIZE - 1 else start + suffixes_per_process
        ranges.append((start, end))
    
    # 启动进程池处理
    with Pool(PROCESS_POOL_SIZE) as pool:
        pool.map(worker_process, ranges)
    
    print("所有数据处理完成!")

额外优化建议

  • 异步请求提速:可以用aiohttp代替requests实现异步API请求和图片下载,IO密集型任务效率会再上一个台阶。
  • 数据库连接池:如果换成PostgreSQL等支持连接池的数据库,可以用连接池代替每个进程创建独立连接,减少连接开销。
  • 监控与日志:添加详细的日志记录(比如用logging模块),方便跟踪处理进度和排查错误。
  • 分布式处理:单台机器不够的话,可以用Celery等分布式框架把任务分发到多台机器并行处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:58:25