Wybierz wzór ładowania i przemieszczania danych za pomocą mssql-python

Sterownik mssql-python oferuje wiele ścieżek zapisu danych do Microsoft SQL. Każda ścieżka nadaje się do różnych obciążeń roboczych. Ten przewodnik pomaga wybrać odpowiedni model na podstawie ilości danych, formatu źródła i semantyki aktualizacji.

Decyduj według obciążenia

Obciążenie Zalecana ścieżka Dlaczego
Załaduj pliki CSV do tabeli Załaduj dane CSV za pomocą kopiowania zbiorczego bulkcopy() Generator obsługuje pliki dowolnego rozmiaru bez ładowania ich do pamięci.
Wstaw pojedynczy wiersz z kodu aplikacji Wstawianie pojedynczego wiersza Niewielki narzut, prosta obsługa błędów, działa z OUTPUT przy zwracaniu wygenerowanych kluczy.
Wstawianie niewielkiej lub średniej partii danych za pomocą kodu aplikacji Wkładki batched Zmniejsza to liczbę podróży w obie strony w porównaniu do pojedynczych insertów.
Załaduj setki wierszy lub więcej z dowolnego źródła Kopiowanie masowe Wkładka TDS bulk to najwydajniejsza droga dla dużych objętości.
Wstaw lub aktualizuj wiersze na podstawie klucza Upsert za pomocą MERGE MERGE obsługuje INSERT, UPDATE, oraz DELETE w jednym zdaniu.
Załaduj DataFrame do tabeli Ładowanie ramek danych bulkcopy_arrow()odczytuje dane Arrow DataFrame bez tworzenia obiektu Python dla każdej wartości.
Załaduj dane Apache Arrow do tabeli Wczytaj dane Arrow bulkcopy_arrow() czyta pamięć Arrow bezpośrednio, bez tworzenia krotek języka Python.
Przygotuj dane za pomocą plików Parquet Scenografia parkietowa Przydatne dla ETL międzysystemowego, gdzie potrzebny jest format pliku pośredni. Parquet to już dane w formacie Arrow, więc ładuje się bez konwersji wierszy.

Załaduj dane CSV za pomocą kopiowania zbiorczego

Ładowanie danych CSV to najczęściej używane pytanie dotyczące pobierania danych w bazie Python. Użyj csv.reader z generatorem zasilającym 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()

Wzorzec generatora utrzymuje stałe zużycie pamięci niezależnie od rozmiaru pliku. W przypadku mapowania kolumn i obsługi tożsamości zobacz operacje kopiowania masowego.

Wkładki jednowierszowe

Używaj pojedynczych insertów do zapisów na poziomie aplikacji, gdzie przetwarzasz jeden rekord na raz. Zastosowanie OUTPUT INSERTED do odzyskania wygenerowanych kluczy:

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

Pojedyncze wkładki są właściwym wyborem, gdy:

  • Na każdą akcję użytkownika (przesłanie formularza, wywołanie API) wstawiasz jeden wiersz.
  • Musisz zweryfikować lub przekształcić każdy wiersz osobno przed wstawieniem.
  • Potrzebujesz natychmiast wstawionego identyfikatora lub innych wygenerowanych wartości.

Wstawianie wsadowe

Używaj executemany(), gdy masz średnią liczbę wierszy i nie potrzebujesz wysokiej wydajności kopiowania zbiorczego:

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() wysyła każdy wiersz jako osobne, parametryzowane zdanie. Gdy przepustowość ma większe znaczenie niż kontrola poszczególnych wierszy, bulkcopy() jest wydajniejsza, ponieważ wykorzystuje protokół zbiorczego wstawiania TDS. Crossover zależy od szerokości wiersza i opóźnienia sieci, ale zazwyczaj występuje w niskich setkach wierszy.

Kopiowanie zbiorcze

Gdy przepustowość ma większe znaczenie niż kontrola na wiersz, używaj bulkcopy(). Wykorzystuje protokół TDS bulk insert, który strumieniuje wiersze zamiast wysyłać jedno zdanie na wiersz:

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

Wskazówki dotyczące wydajności przy kopiowaniu zbiorczym

  • Używaj generatorów do dużych zbiorów danych, aby utrzymać stałe zużycie pamięci.
  • Użyj bulkcopy_arrow(), gdy źródło jest kolumnowe, na przykład DataFrame lub plik Parquet. Pomija konwersję na krotki wierszy w Python.
  • Ustaw batch_size, aby kontrolować, ile wierszy jest wysyłanych w każdej partii TDS. Zacznij od 5 000 i dostosuj w zależności od szerokości wiersza.
  • Używaj blokad tabeli do ładowania wyłącznego: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Wyłącz indeksy przed ładowaniem, a potem buduj ponownie. Ta sekwencja pozwala uniknąć narzutu związanego z utrzymywaniem indeksu podczas ładowania.

W przypadku mapowania kolumn, kolumn tożsamościowych, obsługi NULL i ładowania równoległego, zobacz operacje kopiowania masowego.

Upsert za pomocą MERGE

MERGE to instrukcja Microsoft SQL do warunkowego INSERT, UPDATE i DELETE w ramach jednej operacji. Obsługuje wzorzec "wstaw, jeśli nowy, zaktualizuj, jeśli istnieje", którego programiści Python często potrzebują.

Upsert z pojedynczym rzędem

Dla pojedynczego wiersza użyj MERGE z klauzulą USING, która definiuje aliasy parametrów:

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

Zwiększenie objętości z tablicą etapową

W przypadku zbiorczych operacji upsert najpierw umieść dane w tabeli tymczasowej, a następnie użyj MERGE, aby zaktualizować dane na podstawie tej tabeli. Stosuj insert-or-update jako domyślny wzorzec dla upsertów DataFrame i aktualizacji wsadowych:

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

Ten przykład pokazuje domyślny wzorzec wstawienia lub aktualizacji:

  • INSERT wiersze ze źródła, które nie istnieją w miejscu docelowym (WHEN NOT MATCHED BY TARGET).
  • UPDATE wiersze występujące w obu (WHEN MATCHED).
  • Klauzula OUTPUT raportuje, jakie działania zostały podjęte w każdym wierszu, co jest przydatne w śledzeniach audytu.

Caution

Dodawaj WHEN NOT MATCHED BY SOURCE THEN DELETE tylko wtedy, gdy dane przejściowe stanowią autorytatywną pełną migawkę systemu docelowego. Jeśli partia zawiera tylko zmienione wiersze, ta klauzula usuwa wiersze, które zostały celowo pominięte w źródłowym feedzie.

Jeśli potrzebujesz pełnego uzgodnienia, rozszerz MERGE dopiero po potwierdzeniu, że źródło jest miarodajne dla tabeli docelowej:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

W środowiskach współdzielonych używaj unikalnej globalnej nazwy tabeli tymczasowej dla każdego uruchomienia lub trwałej tabeli pośredniej, aby uniknąć kolizji między jednocześnie uruchamianymi zadaniami.

Kiedy zamiast tego używać oddzielnych instrukcji UPDATE i INSERT

MERGE jest potężny, ale ma przypadki graniczne. Rozważ użycie oddzielnych zdań, gdy:

  • Nie potrzebujesz DELETE logiki. Oddzielne UPDATE, po którym następuje INSERT WHERE NOT EXISTS, jest bardziej czytelne i łatwiejsze do debugowania.
  • Instrukcja MERGE jest na tyle złożona, że trudno przewidzieć działanie mechanizmu blokowania. Oddzielne instrukcje dają wyraźną kontrolę nad szczegółowością blokady.
  • Aktualizujesz tabelę o dużej współbieżności, w której eskalacja blokad MERGE może powodować blokowanie.
# 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()

Załaduj ramki danych

Obiekt DataFrame ma układ kolumnowy, więc załaduj go za pomocą bulkcopy_arrow() zamiast spłaszczać go do krotek wierszy dla bulkcopy().

Dopasuj typy strzałek do kolumn docelowych przed załadowaniem. pyarrow wywnioskowuje float64 dla kolumny numerycznej, której sterownik nie może zmapować na money, decimal ani 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()

Użyj Table.from_pandas() zamiast przekazywać schemat do Table.cast(), który nie może bezpośrednio przekonwertować kolumny typu float na decimal128.

Polars

Polars implementuje interfejs danych Arrow C, więc możesz przekazać sam DataFrame. Odlewaj kolumny najpierw z tego samego powodu:

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

Możesz także ustawić typy podczas odczytu pliku, używając .pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)})

Bezpośrednie przekazanie obiektu DataFrame przekazuje jego bufory do sterownika bez ich kopiowania. df.to_arrow() też działa, ale podczas tej konwersji Polars ponownie koduje kolumny tekstowe, co powoduje skopiowanie wszystkich danych w tych kolumnach.

bulkcopy_arrow() przyjmuje pyarrow.Table, RecordBatch, RecordBatchReader lub dowolny obiekt, który implementuje interfejs danych Arrow C za pośrednictwem __arrow_c_stream__ lub __arrow_c_array__. Przekazanie któregokolwiek z nich do bulkcopy() zgłasza TypeError.

Aby poznać pełne wzorce ładowania DataFrame, zobacz integrację z pandas oraz integrację z Polars.

Załaduj dane Arrow

Gdy dane źródłowe są już w formacie Apache Arrow, cursor.bulkcopy_arrow() ładuje je bez uprzedniego tworzenia krotek w Pythonie.

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

Metoda akceptuje także pyarrow.RecordBatch lub pyarrow.RecordBatchReader, więc możesz przesłać strumieniowo zbiór wyników z cursor.arrow_reader() bezpośrednio do innej tabeli.

Każdy typ kolumny Arrow musi być zgodny z docelowym typem kolumny SQL i zapisujący nie konwertuje między rodzinami typów. Więcej informacji można znaleźć w artykule o integracji z Apache Arrow.

Scenografia parkietowa

Używaj Parquet jako formatu pośredniego podczas migracji danych między systemami lub jeśli potok ETL już generuje pliki Parquet. Plik Parquet wczytuje dane Arrow, więc przekaż go bezpośrednio do bulkcopy_arrow():

import pyarrow.parquet as pq

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

W przypadku dużych plików Parquet iteruj grupy wierszy, aby utrzymać stałe zużycie pamięci. Każda partia to RecordBatch, które bulkcopy_arrow() akceptuje bezpośrednio:

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

Aby przesyłać strumieniowo cały plik w jednym wywołaniu, opakuj partie za pomocą znacznika 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()

Weryfikacja załadowanych danych

Po załadowaniu sprawdź liczbę wierszy i wybiórczo sprawdź dane:

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

W przypadku obciążeń produkcyjnych nie polegaj na transakcji połączenia wywołującego w celu ochrony połączenia bulkcopy() . bulkcopy() otwiera własne wewnętrzne połączenie i niezależnie zatwierdza skopiowane wiersze, więc operacja conn.rollback() na głównym połączeniu nie może ich cofnąć. Dwa podejścia dają atomowość:

  • Ustaw use_internal_transaction=True, aby każda partia była objęta osobną transakcją. Partia, w której podczas przetwarzania wystąpi błąd, jest wycofywana w całości, zamiast pozostawać częściowo załadowana.
  • Aby zweryfikować dane przed ich przeniesieniem, skopiuj je zbiorczo do tabeli przejściowej, zweryfikuj je, a następnie przenieś wiersze do tabeli docelowej, używając elementu INSERT ... SELECT w ramach transakcji na głównym połączeniu. Ponieważ to INSERT działa w Twoim połączeniu, conn.rollback() cofa to, jeśli walidacja się nie powiedzie.
# 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