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

如何在数据查询时跳过报错的数据库表?

问题描述

需求:编写Python代码从源数据库读取表数据并插入目标数据库,遇到schema错误或列错误的表时,跳过该表继续处理下一个表。

问题:移除df=pd.read_sql(select_data_query, conn, chunksize=10000)后,查询部分能正常跳过报错表,但执行插入操作时出现"NoneType object not subscriptable"错误。

原代码:

try:
    source_database=f'{source_database}.dbo'
    cursor = conn.cursor()
    read_table_query="""
        SELECT table_name
        FROM information_schema.tables
        ORDER BY table_name asc;
        """
    cursor.execute(read_table_query)  

    logger.info("Successfully connected to database")

except Exception as e:
    logger.error("Unable to connect to database: %s", str(e))    


for tables in cursor.fetchall():
    tab = tables[0]
    select_data_query = f'select * FROM {source_database}.{tab} where PRCS_DTE > DATEADD(day, -3, CONVERT (date, SYSDATETIME()));'
    df=pd.read_sql(select_data_query, conn, chunksize=10000)

    try:
        cursor.execute(select_data_query)
        print("{} Records selected ".format(cursor.rowcount) + f"from {tab}")

    except Exception as e:
        logger.exception(e)
    

try:
    engine = sa.create_engine(f'mssql+pyodbc://@{target_server}/{target_database}?trusted_connection=yes&driver={target_driver}')

    for chunk_dataframe in df:
        chunk_dataframe.to_sql(f'{tab}', engine, if_exists='append', index=False, method='multi')
        print("{} Records inserted ".format(cursor.rowcount) + f"into {tab}")
        engine.dispose()

except Exception as e:
    logging.exception(e)
问题分析
  • 代码结构混乱:插入逻辑放在表遍历循环之外,导致df和tab仅关联最后一个表,若某表读取失败df为None,遍历df就会触发NoneType错误。
  • Cursor冲突:同时使用conn.cursor()和pd.read_sql,两者共享同一连接的cursor,会导致cursor状态异常,干扰后续操作。
  • 资源释放错误:engine.dispose()放在chunk循环内,会提前关闭数据库连接,后续chunk无法完成插入。
  • 异常捕获范围不合理:读取和插入的异常未在单表处理范围内捕获,无法实现跳过单个报错表的需求。
修正后的代码
import pandas as pd
import sqlalchemy as sa
import logging

# 初始化日志
logger = logging.getLogger(__name__)
logging.basicConfig(level=logging.INFO)

try:
    source_database = f'{source_database}.dbo'
    cursor = conn.cursor()
    read_table_query = """
        SELECT table_name
        FROM information_schema.tables
        ORDER BY table_name asc;
        """
    cursor.execute(read_table_query)  
    logger.info("Successfully connected to source database")
    tables = cursor.fetchall()
    cursor.close()  # 关闭cursor,避免与pd.read_sql冲突
except Exception as e:
    logger.error("Unable to connect to source database: %s", str(e))
    exit()

# 预创建目标数据库引擎
try:
    engine = sa.create_engine(f'mssql+pyodbc://@{target_server}/{target_database}?trusted_connection=yes&driver={target_driver}')
    logger.info("Successfully connected to target database")
except Exception as e:
    logger.error("Unable to connect to target database: %s", str(e))
    exit()

for table_row in tables:
    tab = table_row[0]
    select_data_query = f'select * FROM {source_database}.{tab} where PRCS_DTE > DATEADD(day, -3, CONVERT (date, SYSDATETIME()));'
    
    try:
        # 分批读取数据
        df_chunks = pd.read_sql(select_data_query, conn, chunksize=10000)
        logger.info(f"Start processing table: {tab}")
        
        # 分批插入目标库
        for chunk in df_chunks:
            chunk.to_sql(tab, engine, if_exists='append', index=False, method='multi')
            logger.info(f"Inserted {len(chunk)} records into {tab}")
        
        logger.info(f"Successfully processed table: {tab}")
    
    except Exception as e:
        logger.exception(f"Failed to process table {tab}, skipping it. Error: {str(e)}")
        continue  # 跳过当前表,继续处理下一个

# 所有表处理完成后关闭引擎
engine.dispose()
logger.info("All tables processed, connection closed")
关键修正点
  • 将单表的读取、插入逻辑封装在同一个循环的try-except块内,确保单个表报错时直接跳过,不影响其他表。
  • 提前关闭源库cursor,避免与pd.read_sql的cursor冲突。
  • 预创建目标库引擎,避免重复创建连接,提升处理效率。
  • 将engine.dispose()移到所有表处理完成后,避免提前关闭连接导致插入失败。
  • 用len(chunk)替代cursor.rowcount获取插入记录数,结果更准确且避免cursor状态问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 13:06:01