用 mssql-python 选择数据加载和移动模式

mssql-python驱动程序为将数据写入 Microsoft SQL 提供了多条路径。 每条路径适合不同的工作量。 本指南帮助您根据数据量、源格式和更新语义选择合适的方案。

按工作量决定

工作量 建议的路径 为什么
将CSV文件加载到表格中 使用批量复制加载 CSV 数据 bulkcopy() 使用生成器可以处理任意大小的文件,而无需加载到内存中。
从应用代码中插入一行 单行插入 开销低,错误处理简单直接,可配合 OUTPUT 返回生成的键。
通过应用程序代码插入小到中等规模的批量数据 批量插入 相比单件插入,减少了往返次数。
从任何来源加载数百行甚至更多 批量复制 TDS 批量插入是处理大容量数据最高效的方式。
根据键值插入或更新行 使用 MERGE 执行插入或更新 MERGE 处理 INSERT、 UPDATE和 DELETE ,在一个语句中。
将数据帧加载到表中 加载数据帧 从 pandas 或 Polars 中提取行,并将其传递给 bulkcopy()
通过Parquet文件获取舞台数据 拼花布式舞台 适用于需要中间文件格式的跨系统ETL。

使用批量复制加载 CSV 数据

加载CSV数据是Python数据库工作中最常见的导入问题。 将 csv.reader 与为 bulkcopy() 供电的发电机配合使用:

import csv
import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

# Create a target table
cursor.execute("""
    IF NOT EXISTS (SELECT * FROM sys.tables WHERE name = 'ProductImport')
    CREATE TABLE dbo.ProductImport (
        Name nvarchar(100),
        ProductNumber nvarchar(25),
        ListPrice decimal(10,2)
    )
""")
conn.commit()

def csv_rows(path):
    with open(path, newline="", encoding="utf-8") as f:
        reader = csv.reader(f)
        next(reader)  # Skip header
        for row in reader:
            yield (row[0], row[1], float(row[2]))

result = cursor.bulkcopy(
    "dbo.ProductImport",
    csv_rows("products.csv"),
    batch_size=5000
)
print(f"Loaded {result['rows_copied']} rows")
conn.commit()

生成器模式无论文件大小如何,都能保持内存使用不变。 关于列映射和身份处理,请参见 批量复制操作

单行插入

在应用层写入时使用单次插入,一次处理一条记录。 使用 OUTPUT INSERTED 检索生成的键:

cursor.execute("""
    INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
    OUTPUT INSERTED.Name
    VALUES (%(name)s, %(product_number)s, %(list_price)s)
""", {"name": "Widget", "product_number": "WG-1000", "list_price": 19.99})

inserted_name = cursor.fetchval()
conn.commit()

在以下情况下,单个嵌件是理想选择:

  • 你每进行一次用户操作(表单提交、API 调用),就插入一行。
  • 你需要在插入前逐行验证或转换。
  • 你需要立即获得插入后返回的 ID 或其他生成的值。

批量插入

当行数适中且不需要批量副本吞吐量时使用 executemany()

rows = [
    {"name": "Widget A", "product_number": "WG-1001", "list_price": 19.99},
    {"name": "Widget B", "product_number": "WG-1002", "list_price": 24.99},
    {"name": "Widget C", "product_number": "WG-1003", "list_price": 29.99},
]

cursor.executemany(
    "INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice) VALUES (%(name)s, %(product_number)s, %(list_price)s)",
    rows
)
conn.commit()

executemany() 将每行作为独立的参数化语句发送。 当吞吐量比每行控制更重要时, bulkcopy() 它更高效,因为它采用了TDS批量插入协议。 临界点取决于行宽和网络延迟,但通常在一两百行左右。

批量复制

当吞吐量比每行控制更重要时,使用 bulkcopy()。 它采用TDS批量插入协议,效率远高于逐行插入:

rows = [
    ("Widget A", "WG-1001", 19.99),
    ("Widget B", "WG-1002", 24.99),
    ("Widget C", "WG-1003", 29.99),
]

result = cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
print(f"Loaded {result['rows_copied']} rows")
conn.commit()

批量复制性能提示

  • 使用生成器 处理大数据集,保持内存使用不变。
  • 设置 batch_size 以控制每个 TDS 批次发送的行数。 从 5,000 开始,并根据行宽调整。
  • 使用桌锁 处理独占负载: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True)
  • 加载前先禁用索引,加载后再重建。 这种顺序避免了负载期间的指标维护开销。

关于列映射、身份列、NULL处理和并行加载,请参见 批量复制操作

使用 MERGE 进行插入或更新

MERGE是 Microsoft SQL 中在单个操作中有条件地执行 INSERT、UPDATE 和 DELETE 的语句。 它可处理 Python 开发者常用的“如果是新的则插入,如果已存在则更新”这一模式。

单行插入或更新

对于单行,使用 MERGE,并配合用于定义参数别名的 USING 子句:

cursor.execute("""
    MERGE dbo.ProductImport AS target
    USING (SELECT %(name)s AS Name, %(product_number)s AS ProductNumber, %(list_price)s AS ListPrice) AS source
    ON target.ProductNumber = source.ProductNumber
    WHEN MATCHED THEN
        UPDATE SET
            Name = source.Name,
            ListPrice = source.ListPrice
    WHEN NOT MATCHED THEN
        INSERT (Name, ProductNumber, ListPrice)
        VALUES (source.Name, source.ProductNumber, source.ListPrice);
""", {"name": "Widget A", "product_number": "WG-1001", "list_price": 24.99})
conn.commit()

带备用台的批量上置

对于批量更新,先把数据分阶段放到临时表里,然后再用 MERGE 来更新。 使用插入或更新作为DataFrame上溢和批处理更新的默认模式:

import csv
import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

# Step 1: Create a global temp table for staging
# Note: bulkcopy() requires global temp tables (##), not session temp tables (#)
cursor.execute("""
    IF OBJECT_ID('tempdb..##ProductImportStage') IS NOT NULL
        DROP TABLE ##ProductImportStage;
    CREATE TABLE ##ProductImportStage (
        Name nvarchar(100),
        ProductNumber nvarchar(25),
        ListPrice decimal(10,2)
    )
""")
cursor.commit()

# Step 2: Bulk load into the staging table
def csv_rows(path):
    with open(path, newline="", encoding="utf-8") as f:
        reader = csv.reader(f)
        next(reader)
        for row in reader:
            yield (row[0], row[1], float(row[2]))

cursor.bulkcopy("##ProductImportStage", csv_rows("products_update.csv"), batch_size=5000)

# Step 3: MERGE from staging into the target table
cursor.execute("""
    MERGE dbo.ProductImport AS target
    USING ##ProductImportStage AS source
    ON target.ProductNumber = source.ProductNumber
    WHEN MATCHED THEN
        UPDATE SET
            Name = source.Name,
            ListPrice = source.ListPrice
    WHEN NOT MATCHED BY TARGET THEN
        INSERT (Name, ProductNumber, ListPrice)
        VALUES (source.Name, source.ProductNumber, source.ListPrice)
    OUTPUT $action, INSERTED.ProductNumber, DELETED.ProductNumber;
""")

# Step 4: Read the OUTPUT to see what changed
for row in cursor.fetchall():
    print(f"{row[0]}: inserted={row[1]}, deleted={row[2]}")

conn.commit()

这个示例展示了默认的插入或更新模式:

  • INSERT 源中存在但目标中不存在的行(WHEN NOT MATCHED BY TARGET)。
  • UPDATE 同时存在于两者中的行(WHEN MATCHED)。
  • OUTPUT 子句报告每行所采取的操作,这对审计轨迹非常有用。

注意

仅当暂存数据是目标数据的权威完整快照时,才添加 WHEN NOT MATCHED BY SOURCE THEN DELETE。 如果该批次仅包含发生更改的行,则该子句会删除源馈送中被故意省略的行。

如果你需要完整对账,仅在确认该源是目标表的权威来源后,才扩展 MERGE

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

在共享环境中,请为每次运行使用唯一的全局临时表名称,或使用永久暂存表,以避免并发作业之间发生冲突。

何时使用分开 UPDATE 和 INSERT 语句代替

MERGE 很强大,但也有例外情况。 在以下情况下,请考虑使用单独的语句:

  • 你不需要使用 DELETE 逻辑。 一个单独的 UPDATE,后面跟着 INSERT WHERE NOT EXISTS,可读性更好,也更便于调试。
  • 这个 MERGE 命题足够复杂,锁定行为难以预测。 单独的语句让你明确控制锁的细度。
  • 你正在更新一个高并发的表,其中 MERGE 锁升级可能导致阻塞。
# Simpler alternative: UPDATE then INSERT
cursor.execute("""
    UPDATE dbo.ProductImport
    SET Name = %(name)s, ListPrice = %(list_price)s
    WHERE ProductNumber = %(product_number)s
""", {"name": "Widget A", "list_price": 24.99, "product_number": "WG-1001"})

if cursor.rowcount == 0:
    cursor.execute("""
        INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
        VALUES (%(name)s, %(product_number)s, %(list_price)s)
    """, {"name": "Widget A", "product_number": "WG-1001", "list_price": 24.99})

conn.commit()

加载数据帧

从 pandas 或 Polars 数据框中提取行,并使用 bulkcopy() 加载这些行:

pandas

将 pandas DataFrame 转换为元组并传递给 bulkcopy()

import pandas as pd

df = pd.read_csv("products.csv")

# Convert DataFrame rows to tuples
rows = list(df[["Name", "ProductNumber", "ListPrice"]].itertuples(index=False, name=None))

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

Polars

使用 .rows() 方法将 Polars DataFrame 转换为元组:

import polars as pl

df = pl.read_csv("products.csv")

# Convert Polars DataFrame to list of tuples
rows = df.select(["Name", "ProductNumber", "ListPrice"]).rows()

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

有关 DataFrame 的完整加载模式,请参见 pandas 集成Polars 集成

拼花布式舞台

在系统之间迁移数据时,或当 ETL 流水线已生成 Parquet 文件时,使用 Parquet 作为中间格式:

import pyarrow.parquet as pq

# Read Parquet file
table = pq.read_table("products.parquet")

# Convert to rows for bulkcopy
rows = [tuple(row) for row in zip(*[col.to_pylist() for col in table.columns])]

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

对于大型 Parquet 文件,可按行组读取,以保持内存占用稳定:

import pyarrow.parquet as pq

parquet_file = pq.ParquetFile("products.parquet")

for batch in parquet_file.iter_batches(batch_size=10000):
    rows = [tuple(row) for row in zip(*[col.to_pylist() for col in batch.columns])]
    cursor.bulkcopy("dbo.ProductImport", rows, batch_size=10000)

conn.commit()

验证已加载的数据

加载后,核对行数并抽查数据:

cursor.execute("SELECT COUNT(*) FROM dbo.ProductImport")
count = cursor.fetchval()
print(f"Total rows: {count}")

cursor.execute("""
    SELECT TOP 5 Name, ProductNumber, ListPrice
    FROM dbo.ProductImport
    ORDER BY Name
""")
for row in cursor:
    print(f"  {row.Name} ({row.ProductNumber}): ${row.ListPrice:.2f}")

对于生产负载,不要依赖呼叫连接的事务来保护 bulkcopy() 通话。 bulkcopy() 会打开自己的内部连接,并独立提交复制过来的行,因此主连接上的 conn.rollback() 无法撤销这些行。 两种方法可实现原子性:

  • 设置为 use_internal_transaction=True 将每个批次包裹在自己的事务中。 如果某个批次在处理中途失败,系统会回滚该批次,而不会让其只加载一部分。
  • 为了在将数据提升到正式使用前先进行验证,请先将数据批量复制到暂存表中,验证完成后,再在主连接上的事务中使用 INSERT ... SELECT 将这些行移到目标表。 由于 INSERT 在你的连接上运行,因此如果验证失败,conn.rollback() 就会将其撤销。
# Stage the data. bulkcopy() runs on its own connection, so these rows
# persist regardless of the transaction below.
cursor.bulkcopy("dbo.ProductImport_Stage", rows, batch_size=5000)

try:
    cursor.execute("SELECT COUNT(*) FROM dbo.ProductImport_Stage")
    count = cursor.fetchval()

    if count < expected_count:
        raise ValueError(f"Expected {expected_count} rows, got {count}")

    # This INSERT runs on your connection, so it's covered by the transaction.
    cursor.execute("""
        INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
        SELECT Name, ProductNumber, ListPrice FROM dbo.ProductImport_Stage
    """)
    conn.commit()
except Exception:
    conn.rollback()
    raise