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

使用SQLAlchemy Session批量插入MSSQL时报‘num(INTEGER)非字符串’错误的原因

问题描述

我正在尝试向MSSQL数据库执行批量插入操作:应用从API获取数据并整理为pandas.DataFrame,将“num”列强制转换为int类型(.astype(int))。随后希望通过数据库事务完成以下步骤:

  • 创建临时表
  • 向临时表插入数据
  • 将临时表与生产表合并
  • 无错误时提交事务,出错则回滚

但使用SQLAlchemy Session执行代码时,触发错误**“num (INTEGER) not a string”**;改用with self.engine.begin() as connection方式可正常运行,但需要实现出错时的事务回滚功能。相关代码及错误信息如下:

SQL表结构

CREATE TABLE [dbo].[discountDimensions](
[ID] [int] IDENTITY(1,1) NOT NULL,
[num] [int] NOT NULL,
[name] [varchar](60) NULL,
[posPercent] [float] NULL
)

DataProcessor.py代码

def discounts_dim_to_df(self, data:dict):
    df = pd.DataFrame(data["discounts"])
    df = df.drop(columns=['num', 'name'])
    df['mstrNum'] = df["mstrNum"].astype(int)
    return df.rename(columns={"mstrNum":"num","mstrName":"name"})

Conn.py中Session版本代码(报错)

class Conn:
    def __init__(self, server:str, database:str, username:str, passwd:str): 
        driver = 'ODBC Driver 17 for SQL Server'
        connection_string = f'mssql+pyodbc://{username}:{passwd}@{server}:1433/{database}?driver={driver}'
        try:
            self.engine = create_engine(connection_string, pool_pre_ping=True)
            Session = sessionmaker(bind=self.engine)
            self.session = Session()
        except Exception as err:
            print(f'Failed to connect to {server} -> {database}')
            mm = Mailman()
            mm.connect_to_db_failed(server=server,db=database,error_message=err)

    def update_discounts_db(self, data:pd.DataFrame) -> bool:
        create_temp_sql = '''
        CREATE TABLE #tempDiscountDimensions (
           num INT PRIMARY KEY,
           name VARCHAR(60) NOT NULL,
           posPercent FLOAT NULL
        );
        '''
        merge_sql = '''
        MERGE INTO discountDimensions AS target
        USING #tempDiscountDimensions AS source
        ON target.num = source.num
        WHEN MATCHED THEN
           UPDATE SET target.name = source.name,
                      target.posPercent = source.posPercent
        WHEN NOT MATCHED THEN
           INSERT (num, name, posPercent) VALUES (source.num, source.name, source.posPercent);
        '''  
        try:
            self.session.execute(text("""IF OBJECT_ID('tempdb..#tempDiscountDimensions', 'U') IS NOT NULL
                               DROP TABLE #tempDiscountDimensions;"""))
            self.session.execute(text(create_temp_sql))
            data.to_sql('#tempDiscountDimensions',
                   con=self.session,
                   if_exists='append',
                   index=False,
                   dtype={
                       "num": Integer(),
                       "name": String(),
                       "posPercent": Float()
                   })

            self.session.execute(text(merge_sql))
            self.session.execute(text('DROP TABLE #tempDiscountDimensions'))
            self.session.commit()

        except Exception as err:
            self.session.rollback()
            print(f'Failed to update/insert data from discountDimensions table:\n{err}')
            mm = Mailman()
            mm.process_failed(option='db_connection', error=err)
            return None

错误栈

Traceback (most recent call last):
  File "path\to\folder\main.py", line 81, in <module>
    conn.update_discounts_db(data=data_disc)
  File "path\to\folder\Conn.py", line 82, in update_discounts_db
    data.to_sql('#tempDiscountDimensions',
  File "path\to\folder\_env\Lib\site-packages\pandas\util\_decorators.py", line 333, in wrapper
    return func(*args, **kwargs)
          ^^^^^^^^^^^^^^^^^^^^^
  File "path\to\folder\_env\Lib\site-packages\pandas\core\generic.py", line 3087, in to_sql
    return sql.to_sql(
          ^^^^^^^^^^^
  File "path\to\folder\_env\Lib\site-packages\pandas\io\sql.py", line 842, in to_sql
    return pandas_sql.to_sql(
          ^^^^^^^^^^^^^^^^^^
  File "path\to\folder\_env\Lib\site-packages\pandas\io\sql.py", line 2839, in to_sql
    raise ValueError(f"{col} ({my_type}) not a string")
ValueError: num (INTEGER) not a string

可行的engine.begin()版本代码

try:
    with self.engine.begin() as connection:
        connection.execute(text("""IF OBJECT_ID('tempdb..#tempDiscountDimensions', 'U') IS NOT NULL
                                   DROP TABLE #tempDiscountDimensions;"""))
        connection.execute(text(create_temp_sql))
        data.to_sql('#tempDiscountDimensions',
                con=connection,
                if_exists='append',
                index=False,
                dtype={
                    "num": Integer(),
                    "name": String(),
                    "posPercent": Float()
                })

        connection.execute(text(merge_sql))
        connection.execute(text('DROP TABLE #tempDiscountDimensions'))

问题分析与解决方案

1. Session版本报错原因

使用self.session作为to_sql的con参数时,pandas对Session对象的类型校验逻辑存在偏差,错误地期望num列是字符串类型,而非实际的整数类型,导致抛出类型不匹配错误。

2. 修复Session版本代码(保留事务回滚)

从Session中获取底层连接对象传递给to_sql,即可解决类型校验问题,同时保留Session的事务管理能力:

def update_discounts_db(self, data:pd.DataFrame) -> bool:
    create_temp_sql = '''
    CREATE TABLE #tempDiscountDimensions (
       num INT PRIMARY KEY,
       name VARCHAR(60) NOT NULL,
       posPercent FLOAT NULL
    );
    '''
    merge_sql = '''
    MERGE INTO discountDimensions AS target
    USING #tempDiscountDimensions AS source
    ON target.num = source.num
    WHEN MATCHED THEN
       UPDATE SET target.name = source.name,
                  target.posPercent = source.posPercent
    WHEN NOT MATCHED THEN
       INSERT (num, name, posPercent) VALUES (source.num, source.name, source.posPercent);
    '''  
    try:
        # 从Session获取底层连接
        connection = self.session.connection()
        connection.execute(text("""IF OBJECT_ID('tempdb..#tempDiscountDimensions', 'U') IS NOT NULL
                           DROP TABLE #tempDiscountDimensions;"""))
        connection.execute(text(create_temp_sql))
        data.to_sql('#tempDiscountDimensions',
               con=connection,
               if_exists='append',
               index=False,
               dtype={
                   "num": Integer(),
                   "name": String(),
                   "posPercent": Float()
               })

        self.session.execute(text(merge_sql))
        self.session.execute(text('DROP TABLE #tempDiscountDimensions'))
        self.session.commit()
        return True

    except Exception as err:
        self.session.rollback()
        print(f'Failed to update/insert data from discountDimensions table:\n{err}')
        mm = Mailman()
        mm.process_failed(option='db_connection', error=err)
        return False

3. 关于engine.begin()的事务回滚

with self.engine.begin() as connection本身就是事务性的:上下文管理器会自动在代码块执行成功时提交事务,执行出错时自动回滚,无需额外处理。如果需要手动控制事务,也可以显式开启:

try:
    connection = self.engine.connect()
    trans = connection.begin()
    try:
        connection.execute(text("""IF OBJECT_ID('tempdb..#tempDiscountDimensions', 'U') IS NOT NULL
                                   DROP TABLE #tempDiscountDimensions;"""))
        connection.execute(text(create_temp_sql))
        data.to_sql('#tempDiscountDimensions',
                con=connection,
                if_exists='append',
                index=False,
                dtype={
                    "num": Integer(),
                    "name": String(),
                    "posPercent": Float()
                })

        connection.execute(text(merge_sql))
        connection.execute(text('DROP TABLE #tempDiscountDimensions'))
        trans.commit()
    except Exception as err:
        trans.rollback()
        raise
finally:
    connection.close()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 03:25:54