使用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
相关产品推荐
相关产品推荐

