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

如何使用pymysql实现cursor.execute()的SQL查询分块读取?

实现pymysql分块读取SQL查询结果并拼接为DataFrame

原代码通过pandas的read_sql_query配合chunksize实现分块读取SQL结果,代码如下:

with DATAB.connect().execution_options(autocommit=True) as DB:
    with DB.begin():
        DataContainer = []
        for chunk in pd.read_sql_query("SELECT id, plat_id FROM plats WHERE date_request >= '{}';".format(Previous_Date), con=DB, chunksize=10000):
            DataContainer.append(chunk)
            del chunk
            gc.collect()
        DataFrame = pd.concat(DataContainer, ignore_index=True)
    DB.close()
DATAB.dispose()

现在基于给定的pymysql框架,补充实现相同的分块读取功能,完整代码如下:

import pandas as pd
import gc
import pymysql

Conn = pymysql.connect(user=Config[0], password=Config[1], host=Config[2], database=Config[3], charset='utf8mb4', cursorclass=pymysql.cursors.DictCursor)
chunk_size = 10000
DataContainer = []
try:
    with Conn.cursor() as cursor:
        # 执行参数化查询,避免SQL注入风险
        sql = "SELECT id, plat_id FROM plats WHERE date_request >= %s;"
        cursor.execute(sql, (Previous_Date,))
        
        # 循环分块拉取数据
        while True:
            chunk_data = cursor.fetchmany(chunk_size)
            if not chunk_data:
                break
            # 转换为DataFrame存入容器
            chunk_df = pd.DataFrame(chunk_data)
            DataContainer.append(chunk_df)
            # 清理临时变量释放内存
            del chunk_data, chunk_df
            gc.collect()
        
        # 拼接所有分块数据
        DataFrame = pd.concat(DataContainer, ignore_index=True)
finally:
    # 确保连接最终关闭
    Conn.close()

关键说明:

  • 用cursor.fetchmany(chunk_size)实现分块读取,每次获取指定条数的记录
  • 采用参数化查询(%s占位符)替代字符串格式化,避免SQL注入漏洞
  • 保留原逻辑中的临时变量清理与垃圾回收,控制内存占用
  • 通过try...finally保证数据库连接无论是否异常都会关闭,避免资源泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 13:35:26