Użyj kopiowania masowego z mssql-python

Sterownik mssql-python zawiera funkcję kopiowania masowego, która efektywnie wstawia duże ilości danych do SQL Server, Azure SQL Database, Azure SQL Managed Instance oraz bazy danych SQL w Microsoft Fabric.

Metoda cursor.bulkcopy() oferuje wydajny sposób wczytywania dużych zbiorów danych:

  • Minimalizuje to podróże sieciowe w obie strony.
  • Opcjonalnie omija sprawdzanie ograniczeń podczas obciążenia.
  • Wykorzystuje zoptymalizowany protokół TDS bulk insert.
  • Osiąga przepustowość porównywalną z bcp.exe i SqlBulkCopy.

Natywne rozszerzenie oparte na języku Rust mssql_py_core obsługuje funkcję zbiorczego kopiowania. Działa poza normalnym potokiem kursora execute().

Podstawowy sposób użycia

Wywołaj bulkcopy() na kursorze, przekazując nazwę tabeli docelowej oraz obiekt iterowalny krotek wierszy lub obiektów Row:

Important

Jeśli tworzysz lub zmieniasz docelową tabelę w tej samej sesji, wywołaj conn.commit() przed bulkcopy(). Protokół kopiowania masowego korzysta z osobnego wewnętrznego kanału do odczytu metadanych tabeli, więc niezatwierdzona zmiana DDL może spowodować impas lub timeout.

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']}")

Wartość zwracana

bulkcopy() Zwraca słownik:

Klucz Typ Opis
rows_copied int Liczba pomyślnie skopiowanych wierszy.
batch_count int Liczba przetworzonych partii.
elapsed_time float Czas operacji trwał kilka sekund.

Sygnatura metody

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
)

Mapowania kolumn

Domyślnie kolumny bulkcopy() mapowane są według pozycji porządkowej. Każda kolumna danych odpowiada kolumnie tabeli o tym samym indeksie. Użyj parametru column_mappings , aby nadpisać to zachowanie.

Lista nazw kolumn

Każda pozycja na liście odpowiada indeksowi danych źródłowych:

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

Format zaawansowany: jawne mapowanie indeksów

Każda krotka przyjmuje postać (source_index, target_column_name). Użyj tego formatu, aby pominąć lub zmienić kolejność kolumn:

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

Ładuj z plików

Możesz załadować dane z plików CSV i innych formatów plików, przekazując generator do bulkcopy().

Plik 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")

Duże pliki z grupowaniem

Ustaw batch_size parametr kontrolujący liczbę wierszy, które sterownik wysyła na partię. To podejście sprawdza się dobrze dla dużych plików:

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

Load pandas DataFrames

DataFrame ma układ kolumnowy, więc najszybszą ścieżką jest bulkcopy_arrow(), które wykorzystuje tabelę Arrow, którą pandas już potrafi wygenerować. bulkcopy() przyjmuje krotki reprezentujące wiersze, więc najpierw musisz przekonwertować kolumny na obiekty Pythona.

Przekonwertuj tabelę Arrow do typów kolumn docelowych przed jej załadowaniem. pyarrow wywnioskowuje float64 dla kolumny numerycznej, której sterownik nie może zmapować na money, decimal ani numeric:

import pandas as pd
import pyarrow as pa
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()

target = pa.schema([
    pa.field('ID', pa.int32()),
    pa.field('Name', pa.string()),
    pa.field('Amount', pa.decimal128(19, 4)),   # MONEY
])

table = pa.Table.from_pandas(df, preserve_index=False).cast(target)
result = cursor.bulkcopy_arrow("##PandasDemo", table)

Bez rzutowania ładowanie kończy się niepowodzeniem z komunikatem ValueError: Cannot map Arrow column 'Amount' (Float64) to SQL column 'Amount' (Money). Zbuduj rzutowanie za pomocą Table.cast(), zamiast przekazywać schemat do Table.from_pandas(), które nie może bezpośrednio przekonwertować kolumny float na decimal128. W tym przebiegu wartości NaN przyjmują wartość SQL NULL, więc nie musisz ich najpierw zastępować.

Jeśli potrzebujesz ścieżki wiersz-krotki, itertuples() już daje krotki po przejściu name=None:

data = list(df.itertuples(index=False, name=None))
result = cursor.bulkcopy("##PandasDemo", data)

Załaduj dane Apache Arrow

Używa się cursor.bulkcopy_arrow() do ładowania danych Apache Arrow. Ta metoda odczytuje dane bezpośrednio z pamięci Arrow, więc przed jej wywołaniem nie trzeba tworzyć krotek wierszy w Pythonie.

Argument source przyjmuje pyarrow.Table, pyarrow.RecordBatch, pyarrow.RecordBatchReader lub dowolny obiekt, który udostępnia interfejs danych Arrow C. Pozostałe argumenty są takie same jak bulkcopy().

import mssql_python
import pyarrow as pa

conn = mssql_python.connect(connection_string)

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

cursor.execute("""
    CREATE TABLE ##ArrowDemo (ID INT, Name NVARCHAR(50), Amount FLOAT)
""")

table = pa.table({
    "ID": pa.array([1, 2, 3], type=pa.int32()),
    "Name": pa.array(["Alice", "Bob", "Carol"], type=pa.string()),
    "Amount": pa.array([50000.0, 60000.0, 55000.0], type=pa.float64()),
})

result = cursor.bulkcopy_arrow("##ArrowDemo", table)
print(f"Copied {result['rows_copied']} rows")

Każdy typ kolumny Arrow musi być zgodny z docelowym typem kolumny SQL. Writer nie konwertuje między rodzinami typów, więc przekazanie kolumny float64 do kolumny pieniężnej podnosi ValueError się zanim zostaną zapisane jakiekolwiek wiersze. Używaj decimal128 dla pieniędzy, kolumn dziesiętnych i liczbowych .

Przekazanie źródła strzałki do bulkcopy() podnosi TypeError i kieruje cię do bulkcopy_arrow().

Więcej informacji o wsparciu dla Arrow, w tym o sposobie przesyłania wyników z jednej tabeli do innej, można znaleźć w artykule Integracja z Apache Arrow.

Obsługa wartości NULL

Przekaż None w dowolnej pozycji kolumny, aby wstawić wartość 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)

Kolumny identyfikacyjne

Aby wstawić jawne wartości tożsamościowe, ustaw 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)

Gdy keep_identity=False (ustawienie domyślne), pomiń kolumnę identity z danych i użyj column_mappings, aby wskazać kolumny inne niż identity.

Opcje kopiowania masowego

Parameter Domyślnie Opis
batch_size 0 Liczba wierszy na wsad. 0 pozwala serwerowi wybrać optymalny rozmiar.
timeout 30 Limit czasu operacji w sekundach. Dotyczy samej operacji kopiowania masowego, a nie połączenia wewnętrznego.
keep_identity False Zachowaj wartości identyfikatorów ze źródłowych danych.
check_constraints False Sprawdź ograniczenia tabeli podczas ładowania.
table_lock False Załóż blokadę na poziomie tabeli zamiast blokad na poziomie wierszy.
keep_nulls False Zachowaj wartości NULL zamiast wstawiać domyślne kolumny.
fire_triggers False Spusty ognia INSERT na stole celu.
use_internal_transaction False Opakuj każdą partię w transakcji wewnętrznej.

Note

bulkcopy() otwiera osobne wewnętrzne połączenie z serwerem. To połączenie wewnętrzne dziedziczy limit czasu zapytania kursora: ustaw wartość Connection.timeout na dodatnią przed utworzeniem kursora, a ta sama wartość będzie ograniczać czas próby nawiązania połączenia dla kopiowania zbiorczego. Jeśli limit czasu zapytania kursora wynosi 0, wewnętrzne połączenie korzysta z domyślnego 15-sekundowego limitu czasu nawiązania połączenia. Kursor przyjmuje wartość w momencie utworzenia, więc późniejsza zmiana Connection.timeout nie wpływa na istniejący kursor ani kopię hurtową w trakcie lotu. Zwiększ limit czasu zapytania przed utworzeniem kursora dla wolnych punktów końcowych, punktów końcowych o ograniczonej przepustowości lub o dużych opóźnieniach (na przykład przez VPN lub między regionami).

Zarządzanie błędami

bulkcopy() zgłasza wyjątek, jeśli ładowanie się nie powiedzie, więc opakuj wywołanie w blok try/except, aby przechwycić błędy. Pamiętaj, że bulkcopy() działa na własnym połączeniu wewnętrznym i niezależnie zatwierdza skopiowane wiersze, więc conn.rollback() w głównym połączeniu nie może ich cofnąć. Aby uczynić partię atomową, ustaw use_internal_transaction=True, która owija każdą partię osobną transakcją, która automatycznie cofa się, jeśli partia się nie powiedzie:

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}")

Aby objąć operację ładowania własną logiką walidacji, skopiuj zbiorczo dane do tabeli przejściowej, a następnie przenieś wiersze do tabeli docelowej za pomocą INSERT ... SELECT w ramach transakcji na głównym połączeniu. To INSERT działa na twoim połączeniu, więc conn.rollback() cofa to, jeśli weryfikacja się nie powiedzie.

Authentication

Kopia masowa korzysta z osobnego wewnętrznego kanału, który wymaga własnego tokena. Sterownik automatycznie obsługuje pozyskiwanie tokenów dla obsługiwanych metod uwierzytelniania.

Zarządzana tożsamość (ActiveDirectoryMSI)

Użyj Authentication=ActiveDirectoryMSI dla tożsamości zarządzanej przypisanej przez system lub przez użytkownika. Ta metoda uwierzytelniania jest zalecana dla usług hostowanych w Azure, takich jak maszyny wirtualne Azure, App Service, Functions i 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")

W przypadku tożsamości zarządzanej przypisanej przez użytkownika przekaż identyfikator klienta w parametrach połączenia:

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

Nazwa główna usługi (ActiveDirectoryServicePrincipal)

Zastosowanie Authentication=ActiveDirectoryServicePrincipal do uwierzytelniania przez podmiot usługi (dane klienta).

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

Domyślny łańcuch poświadczeń (ActiveDirectoryDefault)

ActiveDirectoryDefault Testuje kolejno wielu dostawców poświadczeń, takich jak zmienne środowiskowe, tożsamość obciążenia, tożsamość zarządzana i inne. Działa zarówno dla lokalnego rozwoju, jak i usług hostowanych w Azure bez zmian w kodzie.

Więcej informacji o uwierzytelnianiu można znaleźć w artykule Microsoft Entra authentication.

Porady dotyczące wydajności

Poniższe techniki pomagają zmaksymalizować przepustowość kopiowania masowego.

Zacznij od źródła kolumnowego

bulkcopy()Wymaga iterowalnej liczby krotk wierszowych, więc każda wartość musi istnieć jako obiekt Python zanim kopiowanie się rozpocznie. Gdy dane są już w formacie kolumnowym, bulkcopy_arrow() odczytuje bezpośrednio bufory Arrow i pomija ten krok. Pandas lub Polars DataFrame, plik Parquet oraz wynik cursor.arrow() to wszystkie źródła Arrow. Więcej informacji można znaleźć w artykule Load Apache Arrow data.

Używaj generatorów dla dużych zbiorów danych

Generatory minimalizują zużycie pamięci, ponieważ bulkcopy() akceptują dowolne iterowalne:

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

Używaj zamków stołowych do szybszych załadunków

Gdy nie ma współbieżnych odczytów, ustaw table_lock=True, aby zmniejszyć narzut związany z blokowaniem podczas dużego początkowego ładowania.

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

Wyłącz indeksy podczas ładowania

Tymczasowo wyłącz indeksy nieklastrowane przed ładowaniem zbiorczym i odbuduj je później, aby poprawić wydajność:

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

Ładuj tabele równolegle

Otwórz osobne połączenie dla każdej tabeli i uruchamiaj obciążenia równocześnie.

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

Porównanie z alternatywami

Poniższa tabela porównuje kopię masową z innymi metodami wstawiania danych.

Metoda Przypadek użycia Wydajność
cursor.bulkcopy_arrow() Duże zbiory danych, które już są kolumnowe. Najszybszy
cursor.bulkcopy() Duże zbiory danych (ponad 1000 wierszy) pochodzące ze źródeł zorientowanych na wiersze. Szybko
cursor.executemany() Średnie zbiory danych z parametrami. Umiarkowane
cursor.execute() w pętli Małe zbiory danych z prostą logiką. Najwolniejszy