Используйте массовое копирование с mssql-python

Драйвер mssql-python включает функцию массового копирования, которая эффективно вставляет большие объёмы данных в SQL Server, База данных SQL Azure, Управляемый экземпляр SQL Azure и SQL Database в Microsoft Fabric.

Метод cursor.bulkcopy() обеспечивает высокопроизводительный путь для загрузки больших наборов данных:

  • Минимизирует сетевые круговые поездки.
  • Опционально обходит проверку ограничений во время загрузки.
  • Использует оптимизированный протокол TDS bulk insert.
  • Достигает пропускной способности, сопоставимой с bcp.exe и SqlBulkCopy.

Нативное расширение на Rust mssql_py_core обеспечивает работу функции массового копирования. Он работает вне обычного курсорного execute() конвейера.

Основное использование

Вызовите bulkcopy() для курсора, передав имя целевой таблицы и итерируемый объект с кортежами строк или объектами Row:

Important

Если вы создаёте или изменяете целевую таблицу в этом же сеансе, вызовите conn.commit() перед bulkcopy(). Протокол массового копирования использует отдельный внутренний канал для чтения метаданных таблицы, поэтому незафиксированное изменение DDL может привести к тупику или тайм-ауту.

import mssql_python

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

# Create a temp table for the demo
cursor.execute("""
    CREATE TABLE ##BulkDemo (
        ID INT,
        Name NVARCHAR(50),
        Amount MONEY
    )
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

result = cursor.bulkcopy("##BulkDemo", data)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batch(es)")
print(f"Elapsed: {result['elapsed_time']}")

Возвращаемое значение

bulkcopy() Возвращает словарь:

Key Type Описание
rows_copied int Количество успешно скопированных строк.
batch_count int Количество обработанных партий.
elapsed_time float Время выполнения операции в секундах.

Сигнатура метода

cursor.bulkcopy(
    table_name,                    # str – target table (can include schema, e.g. "dbo.MyTable")
    data,                          # Iterable[Tuple | Row] – rows to insert
    batch_size=0,                  # int – rows per batch; 0 = server optimal
    timeout=30,                    # int – operation timeout in seconds
    column_mappings=None,          # List[str] | List[Tuple[int,str]] | None
    keep_identity=False,           # bool – preserve identity values from source
    check_constraints=False,       # bool – check constraints during load
    table_lock=False,              # bool – use table-level lock
    keep_nulls=False,              # bool – preserve NULLs instead of defaults
    fire_triggers=False,           # bool – fire INSERT triggers on target
    use_internal_transaction=False, # bool – use internal transaction per batch
)

Сопоставления столбцов

По умолчанию bulkcopy() сопоставляет столбцы по порядковому номеру. Каждый столбец данных соответствует столбцу таблицы с тем же индексом. Используйте column_mappings параметр, чтобы отменить это поведение.

Список названий колонок

Каждая позиция в списке соответствует индексу исходных данных:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=["ID", "Name", "Amount"],
)

Расширенный формат: явное индексное отображение

Каждый кортеж имеет вид (source_index, target_column_name). Используйте этот формат для пропуска или перестановки столбцов:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=[(0, "ID"), (1, "Name"), (2, "Amount")],
)

Загрузка из файлов

Вы можете загрузить данные из CSV-файлов и других форматов, передав генератор в bulkcopy().

CSV-файл

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("data.csv", "r", ...)
csv_data = """ID,Name,Value
1,Widget,9.99
2,Gadget,24.50
3,Gizmo,4.75
"""

def csv_row_generator(file_obj):
    """Generator that yields tuples from a CSV file object."""
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:  # skip blank lines
            yield (
                int(row[0]),      # ID
                row[1],           # Name
                float(row[2]),    # Value
            )

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

cursor.execute("""
    CREATE TABLE ##CSVImport (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##CSVImport", csv_row_generator(io.StringIO(csv_data)))
print(f"Imported {result['rows_copied']} rows from CSV")

Большие файлы с пакетированием

Настройте batch_size параметр так, чтобы управлять, сколько строк драйвер отправляет за партию. Этот подход хорошо работает для больших файлов:

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("large_file.csv", "r", ...)
csv_data = "\n".join(
    ["ID,Name,Value"] + [f"{i},Item {i},{i * 1.5}" for i in range(1, 201)]
)

def csv_rows(file_obj):
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:
            yield (int(row[0]), row[1], float(row[2]))

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

cursor.execute("""
    CREATE TABLE ##LargeCSV (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy(
    "##LargeCSV",
    csv_rows(io.StringIO(csv_data)),
    batch_size=50,
)
print(f"Imported {result['rows_copied']} rows in {result['batch_count']} batches")

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

Преобразуем pandas DataFrame в список кортежей перед передачей его в bulkcopy():

import pandas as pd
import mssql_python

df = pd.DataFrame({
    'ID': [1, 2, 3],
    'Name': ['Alice', 'Bob', 'Carol'],
    'Amount': [50000.0, 60000.0, 55000.0],
})

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

cursor.execute("""
    CREATE TABLE ##PandasDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [tuple(row) for row in df.itertuples(index=False, name=None)]
result = cursor.bulkcopy("##PandasDemo", data)

Обработка значений NULL

Передайте None в любой позиции столбца, чтобы вставить значение SQL NULL:

cursor.execute("""
    CREATE TABLE ##NullDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", None),       # NULL Amount
    (3, None, 55000.00),    # NULL Name
]

cursor.bulkcopy("##NullDemo", data)

Столбцы идентификаторов

Чтобы вставить явные значения identity, установите keep_identity=True:

cursor.execute("""
    CREATE TABLE ##IdentDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (100, "Alice", 50000.00),
    (200, "Bob", 60000.00),
]

cursor.bulkcopy("##IdentDemo", data, keep_identity=True)

Когда keep_identity=False (по умолчанию) опустите столбец идентичности из данных и используйте column_mappings для таргетирования столбцов без идентичности.

Варианты массового копирования

Parameter По умолчанию Описание
batch_size 0 Ряды на партию. 0 Позволяет серверу выбрать оптимальный размер.
timeout 30 Операция заканчивается через секунды.
keep_identity False Сохраняйте значения идентичности из исходных данных.
check_constraints False Проверьте ограничения таблицы во время загрузки.
table_lock False Установите блокировку на уровне таблицы вместо блокировок на уровне строк.
keep_nulls False Сохраняйте значения NULL вместо вставки значений столбцов по умолчанию.
fire_triggers False Запускайте INSERT триггеры на столе цели.
use_internal_transaction False Оберните каждую партию внутренней транзакцией.

Управление ошибками

bulkcopy() вызывает исключение, если не удаётся загрузить данные, поэтому заключите вызов в блок try/except, чтобы перехватывать ошибки. Имейте в виду, что bulkcopy() использует собственное внутреннее соединение и независимо фиксирует скопированные строки, поэтому conn.rollback() в вашем основном соединении не может отменить эти изменения. Чтобы сделать пакет атомарным, задайте use_internal_transaction=True, который помещает каждый пакет в отдельную транзакцию, автоматически откатываемую при ошибке выполнения пакета:

import mssql_python

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

cursor.execute("""
    CREATE TABLE ##ImportDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

try:
    result = cursor.bulkcopy("##ImportDemo", data, use_internal_transaction=True)
    print(f"Successfully copied {result['rows_copied']} rows")
except (mssql_python.DatabaseError, ValueError) as e:
    # bulkcopy() commits on its own connection, so there's nothing to roll back
    # here. With use_internal_transaction=True, a failed batch is already rolled
    # back on the bulk copy connection.
    print(f"Bulk copy failed: {e}")

Чтобы выполнять загрузку только после собственной проверки, сначала массово скопируйте данные в промежуточную таблицу, а затем перенесите строки в целевую таблицу с помощью INSERT ... SELECT в рамках транзакции в основном соединении. Это INSERT выполняется через ваше соединение, поэтому conn.rollback() отменяет это, если проверка завершается неудачей.

Authentication

Массовое копирование использует отдельный внутренний канал, который требует собственного токена. Драйвер автоматически обрабатывает получение токенов для поддерживаемых методов аутентификации.

Управляемая идентичность (ActiveDirectoryMSI)

Используйте Authentication=ActiveDirectoryMSI для управляемой идентификации, назначаемой системой или пользователем. Этот метод аутентификации рекомендуется для сервисов, размещённых на Azure, таких как виртуальные машины Azure, App Service, Functions и AKS.

import mssql_python

# System-assigned managed identity
conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##MsiDemo (ID INT, Name NVARCHAR(50))")
conn.commit()

result = cursor.bulkcopy("##MsiDemo", [(1, "Alice"), (2, "Bob")])
print(f"Copied {result['rows_copied']} rows")

Для управляемой идентичности, назначаемой пользователем, передайте идентификатор клиента в строке подключения:

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "UID=<client-id>;"
    "Encrypt=yes"
)

Принципал сервиса (ActiveDirectoryServicePrincipal)

Использование Authentication=ActiveDirectoryServicePrincipal для аутентификация основного сервиса (учётные данные клиента).

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryServicePrincipal;"
    "UID=<application-client-id>;"
    "PWD=<client-secret>;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##SpDemo (ID INT, Value FLOAT)")
conn.commit()

result = cursor.bulkcopy("##SpDemo", [(1, 1.5), (2, 2.5)])
print(f"Copied {result['rows_copied']} rows")

Цепочка учетных данных по умолчанию (ActiveDirectoryDefault)

ActiveDirectoryDefault Пробует несколько поставщиков учетных данных последовательно, таких как переменные среды, идентификация рабочей нагрузки, управляемая идентичность и другие. Он работает как для локальной разработки, так и для сервисов, размещённых на Azure, без изменений кода.

Дополнительные сведения об аутентификации см. в разделе Аутентификация Microsoft Entra.

Советы по производительности

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

Используйте генераторы для больших наборов данных

Генераторы минимизируют использование памяти, потому что bulkcopy() принимает любую итерируемую:

def data_generator(count):
    """Generate rows without loading all into memory."""
    for i in range(count):
        yield (i, f"Item {i}", i * 1.5)

cursor = conn.cursor()
cursor.execute("""
    CREATE TABLE ##LargeDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##LargeDemo", data_generator(1000))

Используйте блокировки таблицы для ускорения загрузки данных

Если у вас нет одновременных считывателей, настройте table_lock=True на снижение накладных расходов блокировки при больших первых нагрузках.

result = cursor.bulkcopy(
    "##LargeDemo",
    data,
    table_lock=True,
    batch_size=100000,
)

Отключить индексы во время загрузки

Временно отключите некластерные индексы перед массовой загрузкой и перестройте их после этого для улучшения производительности:

cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##IndexDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
cursor.execute("CREATE NONCLUSTERED INDEX IX_Name ON ##IndexDemo(Name)")
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo DISABLE")
conn.commit()

result = cursor.bulkcopy("##IndexDemo", data)
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo REBUILD")
conn.commit()

Загрузите таблицы параллельно

Откройте отдельное соединение для каждой таблицы и запускайте нагрузки одновременно.

import concurrent.futures

def load_table(table_name, rows):
    conn = mssql_python.connect(connection_string)
    cursor = conn.cursor()
    cursor.execute(f"CREATE TABLE {table_name} (ID INT, Name NVARCHAR(50), Value FLOAT)")
    conn.commit()
    result = cursor.bulkcopy(table_name, rows)
    conn.commit()
    conn.close()
    return result["rows_copied"]

data = [(i, f"Item {i}", i * 1.5) for i in range(100)]

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    futures = [
        executor.submit(load_table, "##Load1", data),
        executor.submit(load_table, "##Load2", data),
        executor.submit(load_table, "##Load3", data),
    ]
    for future in concurrent.futures.as_completed(futures):
        print(f"Loaded {future.result()} rows")

Сравнение с альтернативами

В следующей таблице сравнивается массовое копирование с другими методами вставки данных.

Метод Сценарий использования Производительность
cursor.bulkcopy() Большие наборы данных (более 1000 строк). Самый быстрый
cursor.executemany() Средние наборы данных с параметрами. Умеренно
cursor.execute() в цикле Небольшие наборы данных с простой логикой. Самый медленный