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

Python实现多源SQL并行查询并存储DataFrame的方法问询

并行SQL查询解决方案

核心思路

数据库查询属于IO密集型任务,用多线程(或concurrent.futures.ThreadPoolExecutor)足够提升效率,无需多进程(多进程适合CPU密集型场景,且存在内存隔离问题)。关键是每个任务独立创建数据库连接,避免共享连接引发的线程安全问题,同时通过任务提交机制直接获取返回的DataFrame。

具体实现代码

import pandas as pd
import pyodbc
from hdbcli import dbapi
from concurrent.futures import ThreadPoolExecutor

# 定义查询语句
SSMS_QUERIES = {
    "df1": "SELECT * FROM TABLE1",
    "df2": "SELECT * FROM TABLE2"
}
HANA_QUERIES = {
    "df3": "SELECT * FROM TABLE3"
}

# SSMS查询任务函数:每个任务独立创建连接
def run_ssms_query(query):
    server = "server_name"
    cnxn = pyodbc.connect(f"DRIVER={{SQL Server}};SERVER={server};trusted_connection=Yes")
    df = pd.read_sql(query, cnxn)
    cnxn.close()
    return df

# HANA查询任务函数:每个任务独立创建连接
def run_hana_query(query):
    conn = dbapi.connect(address="your_hana_address", port=30015, user="your_user", password="your_pwd")
    df = pd.read_sql(query, conn)
    conn.close()
    return df

if __name__ == "__main__":
    # 收集所有任务
    tasks = []
    
    # 添加SSMS查询任务
    for df_name, query in SSMS_QUERIES.items():
        tasks.append(("ssms", df_name, query))
    
    # 添加HANA查询任务
    for df_name, query in HANA_QUERIES.items():
        tasks.append(("hana", df_name, query))
    
    # 执行并行任务
    results = {}
    with ThreadPoolExecutor(max_workers=4) as executor:
        # 提交任务并关联结果名称
        future_to_df = {}
        for task_type, df_name, query in tasks:
            if task_type == "ssms":
                future = executor.submit(run_ssms_query, query)
            else:
                future = executor.submit(run_hana_query, query)
            future_to_df[future] = df_name
        
        # 收集返回的DataFrame
        for future, df_name in future_to_df.items():
            results[df_name] = future.result()
    
    # 导出到Excel
    with pd.ExcelWriter("query_results.xlsx") as writer:
        for df_name, df in results.items():
            df.to_excel(writer, sheet_name=df_name, index=False)

关键注意事项

  • 禁止共享数据库连接:数据库连接对象(如pyodbc.connect或dbapi.connect返回的实例)不是线程安全的,必须每个任务独立创建和关闭连接。
  • 线程池大小:max_workers建议设置为查询数量或稍多,避免过多线程导致数据库连接压力过大。
  • 多进程替代方案:如果必须用多进程,将上述代码中的ThreadPoolExecutor替换为ProcessPoolExecutor即可,但注意多进程启动逻辑要放在if __name__ == "__main__"块内,避免重复初始化问题。
  • 结果收集:通过future.result()直接获取每个任务返回的DataFrame,存入字典后即可统一处理导出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 23:15:43