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

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

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

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

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

  • Используйте генераторы для больших наборов данных, чтобы сохранять постоянное использование памяти.
  • 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()

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

Извлечь строки из pandas или Polars DataFrame и загрузить их с помощью 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()

Поларс

Преобразование DataFrame Polars в кортежи с помощью метода .rows():

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.

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

Используйте Parquet как промежуточный формат при миграции данных между системами или когда ваш ETL-конвейер уже создаёт файлы 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