从DB2查询数据写入SQL Server时遇字符截断错误求助
问题描述
运行Python多线程脚本从IBM DB2拉取数据写入SQL Server Express时,大量表和列出现字符截断错误:
Query Error on TABLE1: ('22001', '[22001] [IBM][System i Access ODBC Driver]Column 1: COLUMN1 - Input data is too big to fit into field (30207) (SQLExecDirectW); [22001] [IBM][System i Access ODBC Driver]Column 1: Character data right truncation. (30126)')
已尝试调整SQL Server字段类型(从NVARCHAR(MAX)改为NVARCHAR(255)/NVARCHAR(100))、使用TRIM去除空格,但问题未解决。
相关代码
Python脚本
import threading import time import pyodbc class DB2Worker(threading.Thread): def __init__(self, interval, query, params, target_table, db2_conn_str, local_conn_str): super().__init__() self.interval = interval self.query = query self.params = params self.target_table = target_table self.db2_conn_str = db2_conn_str self.local_conn_str = local_conn_str self.daemon = True self.db2_conn = None def _get_db2_connection(self): """Maintains the connection to the remote DB2 source.""" while self.db2_conn is None: try: self.db2_conn = pyodbc.connect(self.db2_conn_str) self.db2_conn.autocommit = True except Exception as e: print(f"DB2 Connection Error: {e}. Retrying in 5s...") time.sleep(5) def _save_to_local_sql(self, data, db2_description): """Dynamic mapping to Local SQL Express.""" if not data: return try: column_names = [col[0] for col in db2_description] cols_str = ", ".join(column_names) placeholders = ",".join(["?" for _ in column_names]) insert_sql = f"INSERT INTO {self.target_table} ({cols_str}) VALUES ({placeholders})" with pyodbc.connect(self.local_conn_str) as local_conn: with local_conn.cursor() as cursor: cursor.fast_executemany = True cursor.executemany(insert_sql, data) local_conn.commit() print(f"[{self.target_table}] Successfully synced {len(data)} rows.") except Exception as e: print(f"Local SQL Save Error for {self.target_table}: {e}") def run(self): self._get_db2_connection() while True: try: cursor = self.db2_conn.cursor() if self.params: cursor.execute(self.query, self.params) else: cursor.execute(self.query) if cursor.description: results = cursor.fetchall() if results: self._save_to_local_sql(results, cursor.description) cursor.close() except (pyodbc.Error, Exception) as e: print(f"Query Error on {self.target_table}: {e}") self.db2_conn = None self._get_db2_connection() time.sleep(self.interval) DB2_CONNECTION_STRING = ( "DRIVER={iSeries Access ODBC Driver};" "DATABASE=DB2;" "SYSTEM=192.168.10.1;" "PORT=446;" "PROTOCOL=TCPIP;" "UID=USERNAME;" "PWD=PASSWORD;" "ConnectTimeout=10;" "LongDataCompat=1;" "TrimChar=1;" ) LOCAL_SQL_CONNECTION_STRING = ( "Driver={ODBC Driver 18 for SQL Server};" "Server=SERVER\\SERVER;" "Database=DB;" "Trusted_Connection=yes;" "Encrypt=yes;" "TrustServerCertificate=yes;" ) query_tasks = [ # 1. Queries WITH tags (using '?' and raw strings) (1, 'SELECT "G1ST" FROM "DB2"."TABLE1" WHERE "G1R" = ?', (r'192.168.64.1\Tag1',), "TABLE1"), (0.5, '''SELECT "FWF","FWCF","FWI","FWWSF3","FWWSI3" FROM "DB2"."TABLE2" WHERE "FWM" = 'LOC1' AND "FWR" = ? ORDER BY "FWF" ASC''', (r'192.168.64.1\Tag2',), "TABLE2"), (1, 'SELECT "G1SB" FROM "DB2"."TABLE1" WHERE "G1R" = ?', (r'192.168.64.1\Tag3',), "TABLE1"), # 2. Queries WITHOUT tags (Standard SQL) (1.0, '''SELECT "B1R","B1B" FROM "DB2"."TABLE3" WHERE "B1W" = 'LOC3' ORDER BY "B1X" DESC , "B1T" DESC''', None, "TABLE3"), (1.0, '''SELECT "GSR","GSS" FROM "DB2"."TABLE4" WHERE "GSM" = 'LOC1' AND "GSA" < 7 ORDER BY "GSS" ASC''', None, "TABLE4"), (1.0, '''SELECT "GSR" FROM "DB2"."TABLE4" WHERE "GSM" = 'LOC1' AND "GSA" = 5 ORDER BY "GSS" DESC''', None, "TABLE4"), (0.5, '''SELECT "G7G","G7F" FROM "DB2"."TABLE5" WHERE "G7CM" = 'LOC1' ORDER BY "G7CJ" DESC , "G7T" DESC''', None, "TABLE5"), (1.5, '''SELECT "G7W","G7R","G7MAP" FROM "DB2"."TABLE6" WHERE "G7MAC" = 'LOC1' ''', None, "TABLE6"), (1.5, '''SELECT "G7W","G7R","G7MAP" FROM "DB2"."TABLE6" WHERE "G7MAC" = 'LOC1' ''', None, "TABLE6") ] if __name__ == "__main__": for interval, sql, params, target_table in query_tasks: worker = DB2Worker(interval, sql, params, target_table, DB2_CONNECTION_STRING, LOCAL_SQL_CONNECTION_STRING) worker.start() time.sleep(0.5) while True: time.sleep(10)
SQL Server建表语句
USE DB; GO IF OBJECT_ID('TABLE3', 'U') IS NOT NULL DROP TABLE TABLE3; CREATE TABLE TABLE3 ( B1R NVARCHAR(100), B1B NVARCHAR(100), B1W NVARCHAR(100), B1X INT, B1T INT ); IF OBJECT_ID('TABLE1', 'U') IS NOT NULL DROP TABLE TABLE1; CREATE TABLE TABLE1 ( G1R NVARCHAR(100), G1ST NVARCHAR(100), G1SB NVARCHAR(100), G1AF FLOAT, G1AI FLOAT, G1AT FLOAT, G1AD FLOAT, G1CL FLOAT, G1DTYP NVARCHAR(100), G1TDTE INT, G1TMCH NVARCHAR(100), G1DDTE INT, G1DMCH NVARCHAR(100), G1DLOT NVARCHAR(100), G1TWGT FLOAT, G1PRDG NVARCHAR(100), G1SIZE NVARCHAR(100), G1CM NVARCHAR(255) ); IF OBJECT_ID('TABLE7', 'U') IS NOT NULL DROP TABLE TABLE7; CREATE TABLE TABLE7 ( FYS NVARCHAR(100), FYC NVARCHAR(100) ); IF OBJECT_ID('TABLE8', 'U') IS NOT NULL DROP TABLE TABLE8; CREATE TABLE TABLE8 ( FZST NVARCHAR(100), FZSI NVARCHAR(100), FZB NVARCHAR(100), FZR INT, FZG FLOAT, FZTW FLOAT, FZTS FLOAT, FZTH FLOAT, FZT FLOAT, FZTL FLOAT, FZSTE INT, FZPAT NVARCHAR(100), FZL FLOAT, FZQTY1 FLOAT, FZL2 FLOAT, FZY FLOAT, FZI INT ); IF OBJECT_ID('TABLE9', 'U') IS NOT NULL DROP TABLE TABLE9; CREATE TABLE TABLE9 ( FZS NVARCHAR(100), FZD NVARCHAR(100), FZF NVARCHAR(100) ); IF OBJECT_ID('TABLE10', 'U') IS NOT NULL DROP TABLE TABLE10; CREATE TABLE TABLE10 ( FYM NVARCHAR(100), FYS NVARCHAR(100), FYO FLOAT ); IF OBJECT_ID('TABLE11', 'U') IS NOT NULL DROP TABLE TABLE11; CREATE TABLE TABLE11 ( G51S NVARCHAR(100), G51I NVARCHAR(100) ); IF OBJECT_ID('TABLE4', 'U') IS NOT NULL DROP TABLE TABLE4; CREATE TABLE TABLE4 ( GSM NVARCHAR(100), GSR NVARCHAR(100), GSS INT, GSA INT ); IF OBJECT_ID('TABLE5', 'U') IS NOT NULL DROP TABLE TABLE5; CREATE TABLE TABLE5 ( G7G NVARCHAR(100), G7F NVARCHAR(100), G7CM NVARCHAR(100), G7CJ INT, G7T INT ); IF OBJECT_ID('TABLE2', 'U') IS NOT NULL DROP TABLE TABLE2; CREATE TABLE TABLE2 ( FWM NVARCHAR(100), FWR NVARCHAR(100), FWF INT, FWCF FLOAT, FWI FLOAT, FWWSF3 FLOAT, FWWSI3 FLOAT ); IF OBJECT_ID('TABLE6', 'U') IS NOT NULL DROP TABLE TABLE6; CREATE TABLE TABLE6 ( G7MAC NVARCHAR(100), G7W NVARCHAR(100), G7R FLOAT, G7MAP FLOAT );
解决方案
1. 确认DB2源字段实际长度
先查询DB2中对应表的字段定义,确保SQL Server的字段长度不小于DB2的字段长度:
-- 在DB2中执行 SELECT COLUMN_NAME, DATA_TYPE, CHARACTER_MAXIMUM_LENGTH FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = 'TABLE1' AND TABLE_SCHEMA = 'DB2';
比如如果DB2中G1ST是VARCHAR(200),SQL Server中定义为NVARCHAR(100)就会导致截断,需要调整为NVARCHAR(200)或更大。
2. 修改DB2 ODBC连接参数
当前连接字符串中TrimChar=1可能未生效,尝试添加字符集参数或调整长数据兼容设置:
DB2_CONNECTION_STRING = ( "DRIVER={iSeries Access ODBC Driver};" "DATABASE=DB2;" "SYSTEM=192.168.10.1;" "PORT=446;" "PROTOCOL=TCPIP;" "UID=USERNAME;" "PWD=PASSWORD;" "ConnectTimeout=10;" "LongDataCompat=0;" # 尝试切换为0 "TrimChar=1;" "CHARSET=UTF-8;" )
3. 在DB2查询中显式截断字段
如果无法调整SQL Server字段长度,在DB2查询时直接截断过长字符,确保匹配目标字段长度:
SELECT LEFT("G1ST", 100) AS "G1ST" FROM "DB2"."TABLE1" WHERE "G1R" = ?
4. 在Python中预处理数据
过滤特殊字符并截断到目标字段长度,避免隐性字节溢出:
def _process_data(self, data, target_length=100): processed = [] for row in data: new_row = [] for val in row: if isinstance(val, str): # 截断到目标长度 val = val[:target_length] # 移除不可见控制字符 val = ''.join(c for c in val if ord(c) >= 32) new_row.append(val) processed.append(tuple(new_row)) return processed # 在_save_to_local_sql前调用处理数据 processed_data = self._process_data(results) self._save_to_local_sql(processed_data, db2_description)
5. 关闭fast_executemany测试
cursor.fast_executemany = True可能存在类型映射问题,临时关闭验证是否解决:
# 注释掉该行 # cursor.fast_executemany = True
内容的提问来源于stack exchange,提问作者user32732426

