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

从DB2查询数据写入SQL Server时遇字符截断错误求助

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.02 01:47:27