Выберите шаблон загрузки и перемещения данных с помощью mssql-python

Драйвер mssql-python предоставляет несколько путей для записи данных в Microsoft SQL. Каждый путь подходит под разные нагрузки. Это руководство поможет вам выбрать правильный вариант, исходя из объёма данных, формата исходного кода и семантики обновления.

Решайте по нагрузке

Рабочая нагрузка Рекомендуемый путь Почему
Загрузите CSV-файлы в таблицу Загрузите данные CSV с помощью пакетного копирования bulkcopy() Генератор обрабатывает файлы любого размера без загрузки их в память.
Вставьте одну строку из кода приложения Однорядные вставки Низкие накладные расходы, простая обработка ошибок, работает с OUTPUT для возврата сгенерированных ключей.
Вставьте небольшой или средний пакет данных из кода приложения Пакетные вставки Это уменьшает количество круговых поездок по сравнению с одиночными вставками.
Загрузите сотни и более строк из любого источника Массовое копирование TDS bulk insert — самый эффективный способ при больших объёмах.
Вставляйте или обновляйте строки на основе ключа Апсерт с MERGE MERGE обрабатывает INSERT, UPDATE и DELETE в одной инструкции.
Загрузите DataFrame в таблицу Загрузить фреймы данных bulkcopy_arrow() читает данные Arrow в DataFrame, не создавая объект Python для каждого значения.
Загрузите данные Apache Arrow в таблицу Загрузить данные Arrow bulkcopy_arrow() читает данные из памяти Arrow напрямую, без создания кортежей Python.
Промежуточное хранение данных в файлах Parquet Инсценировка паркета Полезно для межсистемного ETL, где требуется промежуточный формат файла. Parquet уже содержит данные Arrow, поэтому загружается без преобразования строк.

Загрузите данные 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, который потоково передаёт строки вместо отправки отдельной SQL-инструкции для каждой строки:

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()

Советы по производительности для массового копирования

  • Используйте генераторы для больших наборов данных, чтобы сохранять постоянное использование памяти.
  • Используйте bulkcopy_arrow(), когда источник имеет столбчатую структуру, например DataFrame или файл Parquet. Преобразование в кортежи строк Python пропускается.
  • Set batch_size задает, сколько строк отправлять в каждом пакете TDS. Начни с 5 000 и корректируй по ширине рядов.
  • Используйте блокировки таблицы для монопольной загрузки: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Отключите индексы перед загрузкой, а потом перестройте после. Эта последовательность позволяет избежать накладных расходов по поддержанию индекса во время нагрузки.

Сведения о сопоставлениях столбцов, столбцах идентификаторов, обработке NULL и параллельной загрузке см. в разделе Операции массового копирования.

Обновление или вставка с MERGE

MERGE — это инструкция SQL в Microsoft, которая позволяет условно выполнять 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 используйте для обновления из неё. Используйте шаблон insert-or-update в качестве шаблона по умолчанию для обновлений 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()

Загрузить таблицы данных

DataFrame имеет столбчатую структуру, поэтому загружайте его с помощью bulkcopy_arrow(), а не преобразовывайте в кортежи строк для bulkcopy().

Сопоставьте типы стрелок с столбцами назначения перед загрузкой. pyarrow определяет float64 для числового столбца, который драйвер не может сопоставить с money, decimal или numeric.

pandas

import pandas as pd
import pyarrow as pa

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

target = pa.schema([
    pa.field("Name", pa.string()),
    pa.field("ProductNumber", pa.string()),
    pa.field("ListPrice", pa.decimal128(19, 4)),   # MONEY
])

table = pa.Table.from_pandas(
    df[["Name", "ProductNumber", "ListPrice"]], preserve_index=False
).cast(target)

cursor.bulkcopy_arrow("dbo.ProductImport", table)
conn.commit()

Используй Table.cast() , а не передавай схему в Table.from_pandas(), которая не может напрямую преобразовать плавающий столбец в decimal128 .

Поларс

Polars реализует интерфейс данных Arrow C, так что вы можете передавать сам DataFrame. Сначала отливайте колонки по той же причине:

import polars as pl

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

cursor.bulkcopy_arrow(
    "dbo.ProductImport",
    df.select([
        "Name",
        "ProductNumber",
        pl.col("ListPrice").cast(pl.Decimal(19, 4)),   # MONEY
    ]),
)
conn.commit()

Вы также можете задать типы при чтении файла, используя pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

При прямой передаче DataFrame его буферы передаются драйверу без создания копии. df.to_arrow() тоже работает, но при этом преобразовании Polars перекодирует строковые столбцы, что приводит к копированию всех строковых данных.

bulkcopy_arrow() принимает pyarrow.Table, a RecordBatch, a RecordBatchReaderили любой объект, реализующий интерфейс данных Arrow C через __arrow_c_stream__ или __arrow_c_array__. Передача любого из них в bulkcopy() вызывает TypeError.

Все варианты загрузки DataFrame см. в разделах интеграция с pandas и интеграция с Polars.

Загрузить данные Arrow

Когда исходные данные уже находятся в формате Apache Arrow, cursor.bulkcopy_arrow() загружает их без предварительного создания кортежей в Python.

from decimal import Decimal

import pyarrow as pa

# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()

table = pa.table({
    "Name": pa.array(["Widget", "Gadget"], type=pa.string()),
    "ProductNumber": pa.array(["WI-1000", "GA-2000"], type=pa.string()),
    "ListPrice": pa.array([Decimal("29.99"), Decimal("49.99")], type=pa.decimal128(10, 2)),
})

result = cursor.bulkcopy_arrow("dbo.ProductImport", table, batch_size=5000)
print(f"Copied {result['rows_copied']} rows")

Метод также принимает pyarrow.RecordBatch или pyarrow.RecordBatchReader, поэтому вы можете напрямую передавать набор результатов из cursor.arrow_reader() в другую таблицу.

Каждый тип столбца Arrow должен быть совместим с соответствующим ему типом столбца SQL, а модуль записи не преобразует данные между семействами типов. Для получения дополнительной информации см. раздел интеграция Apache Arrow.

Инсценировка паркета

Используйте Parquet как промежуточный формат при миграции данных между системами или когда ваш ETL-конвейер уже создаёт файлы Parquet. Файл Parquet считывается как данные Arrow, поэтому передайте его напрямую в bulkcopy_arrow():

import pyarrow.parquet as pq

cursor.bulkcopy_arrow("dbo.ProductImport", pq.read_table("products.parquet"))
conn.commit()

Для больших файлов Parquet последовательно обрабатывайте группы строк, чтобы потребление памяти оставалось постоянным. Каждая партия — это RecordBatch, которую bulkcopy_arrow() принимает напрямую:

import pyarrow.parquet as pq

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

for batch in parquet_file.iter_batches(batch_size=10000):
    cursor.bulkcopy_arrow("dbo.ProductImport", batch)

conn.commit()

Чтобы передать весь файл потоком за один вызов, оберните пакеты в тег RecordBatchReader:

import pyarrow as pa
import pyarrow.parquet as pq

parquet_file = pq.ParquetFile("products.parquet")
reader = pa.RecordBatchReader.from_batches(
    parquet_file.schema_arrow, parquet_file.iter_batches(batch_size=10000)
)

cursor.bulkcopy_arrow("dbo.ProductImport", reader)
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