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

如何让Python程序无需等待响应,并发向数据库发送请求提速数据匹配?

优化百万级数据的数据库匹配程序(IO密集型任务提速)

问题核心分析

当前程序单线程串行执行远程数据库查询,大部分时间都在等待数据库IO响应,导致CPU利用率极低、处理速度慢。针对IO密集型场景,多线程并发执行查询是最优方案——利用等待数据库响应的时间处理其他记录,同时避免多进程带来的高开销(Python GIL在IO等待时会自动释放,不影响多线程并发效率)。

优化方案与代码实现

关键优化点:

  • 用ThreadPoolExecutor实现多线程并发
  • 每个线程独立创建数据库连接(数据库连接通常非线程安全,禁止共用)
  • 对CSV写入操作加锁,避免多线程写入冲突
  • 修复原代码中重复调用fetch_assoc、SQL语句缺失FROM的问题
import ibm_db
import sys
from concurrent.futures import ThreadPoolExecutor, as_completed
import threading

# 全局锁:保证CSV写入操作线程安全
csv_lock = threading.Lock()

def print_row_to_csv_safe(row):
    """线程安全的CSV写入包装函数"""
    with csv_lock:
        print_row_to_csv(row)  # 调用原有的CSV写入函数

def process_single_record(record):
    """单条记录的查询与处理逻辑,每个线程独立执行"""
    # 每个线程单独创建数据库连接
    db_conn = ibm_db.connect("你的数据库连接字符串", "", "")
    if not db_conn:
        print(f"数据库连接失败: {ibm_db.conn_errormsg()}")
        return

    try:
        # 修复SQL语句的FROM缺失问题
        stmt = ibm_db.prepare(
            db_conn,
            "select trim(Name) || ' ' || trim(FTHFNAME) || ' ' || trim(FTHSNAME) || ' ' || trim(FTHTNAME) AS NAME, BIRTHDATE, IDNUMSTAT FROM PEOPLE where NAME = ? with ur;",
            {ibm_db.SQL_ATTR_CURSOR_TYPE: 3}
        )
        ibm_db.bind_param(stmt, 1, record['NAME'])
        ibm_db.execute(stmt)

        # 仅调用一次fetch_assoc,避免重复获取导致数据丢失
        query_result = ibm_db.fetch_assoc(stmt)
        if query_result:
            print("Matched")
            print_row_to_csv_safe(query_result)
        else:
            print("Not matched")
            print_row_to_csv_safe({"Not Matched": "Not found"})
    except Exception as e:
        print(f"处理记录失败[{record['NAME']}]: {ibm_db.stmt_errormsg(stmt)}")
    finally:
        # 确保线程结束时关闭数据库连接
        ibm_db.close(db_conn)

def run_query_parallel(match_data, max_workers=15):
    """并发执行批量记录匹配"""
    # max_workers需根据数据库允许的最大并发连接数调整,建议10-20
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        # 提交所有记录的处理任务
        futures = [executor.submit(process_single_record, record) for record in match_data]
        
        # 等待所有任务完成,可在此处集中处理异常
        for future in as_completed(futures):
            try:
                future.result()
            except Exception as e:
                print(f"任务执行异常: {str(e)}")

额外优化建议

  1. 数据库连接池:如果创建连接开销较大,可使用ibm_db.pool实现连接池复用,减少连接创建销毁的开销
  2. 批量查询替代单条查询:若业务允许,将多个NAME合并为IN查询(比如一次查50-100条),能大幅减少数据库请求次数,比多线程优化效果更显著
  3. 并发数控制:根据数据库的最大并发连接数限制调整max_workers,避免压垮数据库

内容的提问来源于stack exchange,提问作者Mahmoud Khaled Sayed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 00:22:08