如何在数据查询时跳过报错的数据库表?
问题描述
需求:编写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
相关产品推荐
相关产品推荐

