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

如何在Python代码中检查列是否存在以规避SQL列不存在错误

问题与解决方案

原始代码

import pyodbc
import logging
import json
import pandas as pd
import sqlalchemy as sa
import warnings

def read_query():
    logger = logging.getLogger()
    logger.setLevel(logging.INFO)

    with open(r'upsert_config.json','r') as ts:
        config = json.load(ts)

    source_driver = config['source_database_params']['source_driver']
    source_server = config['source_database_params']['source_server']
    source_database = config['source_database_params']['source_database']
    target_driver = config['target_database_params']['target_driver']
    target_server = config ['target_database_params']['target_server']
    target_database = config['target_database_params']['target_database']


    conn = pyodbc.connect(
    f'Driver={source_driver};'
    f'Server={source_server};'
    f'Database={source_database};'
    f'Driver={target_driver};'
    f'Server={target_server};'
    f'Database={target_database};'
    f'Trusted_Connection=yes;'
    f'MARS_Connection=Yes'
    )

    try:
        source_database=f'{source_database}.dbo'
        cursor = conn.cursor()
        read_table_query="""        
        SELECT distinct table_name
        FROM information_schema.columns
        WHERE COLUMN_NAME in ('PRCS_DTE', 'EFF_DTE', 'PRCS_RUN_DTE', 'reportDate', 'DateofData', 'InsertDate')
        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_query1 = f'SELECT * FROM {source_database}.{tab} WHERE PRCS_DTE > DATEADD(day, -2, CONVERT(date, SYSDATETIME()));'
        select_data_query2 = f'SELECT * FROM {source_database}.{tab} WHERE EFF_DTE > DATEADD(day, -2, CONVERT(date, SYSDATETIME()));'
        try:
            df=pd.read_sql(select_data_query1, conn, chunksize=10000)
            df2=pd.read_sql(select_data_query2, conn, chunksize=10000)
            warnings.filterwarnings("ignore")
        except Exception as e:
            logger.exception(e)
            continue       
        
        try:
            engine = sa.create_engine(f'mssql+pyodbc://@{target_server}/{target_database}?trusted_connection=yes&driver={target_driver}')

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

        except Exception as e:
            logging.exception(e)


           
read_query()

遇到的问题

筛选出包含PRCS_DTE、EFF_DTE等指定日期列的表后,遍历表时使用固定列名编写查询语句,导致部分表因不存在对应列触发错误:

  • Invalid column name 'PRCS_DTE'
  • Invalid column name 'EFF_DTE'
  • Invalid column name 'PRCS_RUN_DTE'等

尝试新增多个查询语句后,仍会在无对应列的表上报错,仅存在对应列的表能正常运行,纠结是用if...else...判断列是否存在,还是直接忽略错误。

解决方案

优先选择动态判断列存在性并生成查询,而非单纯忽略错误——忽略错误可能会遗漏本应提取的数据,动态判断则能精准处理每个表的情况:

修改思路

  1. 遍历每个表时,先查询该表实际存在哪些目标日期列
  2. 针对每个存在的日期列,动态生成对应的查询语句,提取近2天数据
  3. 统一处理提取到的数据,插入目标表

修改后的核心代码

# 替换原有的for tables循环部分
target_date_columns = ['PRCS_DTE', 'EFF_DTE', 'PRCS_RUN_DTE', 'reportDate', 'DateofData', 'InsertDate']

for tables in cursor.fetchall():
    tab = tables[0]
    # 查询当前表存在的目标日期列
    check_columns_query = f"""
    SELECT COLUMN_NAME 
    FROM information_schema.columns 
    WHERE TABLE_NAME = '{tab}' 
    AND COLUMN_NAME IN ({', '.join([f"'{col}'" for col in target_date_columns])})
    """
    cursor.execute(check_columns_query)
    existing_cols = [row[0] for row in cursor.fetchall()]
    
    if not existing_cols:
        logger.info(f"Table {tab} has no target date columns, skipping")
        continue
    
    engine = sa.create_engine(f'mssql+pyodbc://@{target_server}/{target_database}?trusted_connection=yes&driver={target_driver}')
    
    for col in existing_cols:
        # 动态生成查询语句
        select_query = f'SELECT * FROM {source_database}.{tab} WHERE {col} > DATEADD(day, -2, CONVERT(date, SYSDATETIME()));'
        try:
            df = pd.read_sql(select_query, conn, chunksize=10000)
            for chunk_dataframe in df:
                rowcount = chunk_dataframe.to_sql(f'{tab}', engine, if_exists='append', index=False, method='multi', chunksize=10000)
                logger.info(f"{rowcount} Records inserted into {tab} via column {col}")
        except Exception as e:
            logger.exception(f"Failed to process table {tab} with column {col}: {str(e)}")
            continue
    
    engine.dispose()

优势

  • 精准匹配每个表的实际列,避免无意义的报错
  • 不会遗漏任何包含目标日期列的表的数据提取
  • 逻辑清晰,便于后续扩展更多日期列

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 02:15:58