Vyberte vzor načítání a pohybu dat pomocí mssql-python

Ovladač mssql-python poskytuje více cest pro zápis dat do Microsoft SQL. Každá cesta odpovídá jiným pracovním zátěžím. Tento průvodce vám pomůže vybrat ten správný na základě objemu dat, formátu zdroje a aktualizace sémantiky.

Rozhodujte podle pracovní zátěže

Pracovní zátěž Doporučená cesta Proč
Načtení CSV souborů do tabulky Načtení dat CSV pomocí hromadného kopírování bulkcopy() pomocí generátoru zpracovává soubory libovolné velikosti, aniž by je bylo nutné načítat do paměti.
Vložte jeden řádek z aplikačního kódu Vkládání jednotlivých řádků Nízká režie, přímé řešení chyb, funguje s OUTPUT pro vrácení generovaných klíčů.
Vložte malou až střední dávku z aplikačního kódu Dávkové vložky Snižuje to počet cest tam a zpět ve srovnání s jednotlivými vložkami.
Načtěte stovky řádků nebo více z jakéhokoliv zdroje hromadné kopírování Hromadné vkládání přes TDS je nejefektivnější způsob pro velké objemy dat.
Vkládat nebo aktualizovat řádky na základě klíče Upsert pomocí MERGE MERGE pracuje s INSERT, UPDATE a DELETE v jednom příkazu.
Načíst DataFrame do tabulky Načtení DataFrameů Extrahujte řádky z pandas nebo Polars a předejte je do bulkcopy().
Data o fázi prostřednictvím souborů Parquet Parketové uspořádání Užitečné pro ETL napříč systémy, kde je potřeba mezilehlé formátování souboru.

Načíst data CSV pomocí nástroje Bulk Copy

Načítání CSV dat je nejčastější otázka při práci s databázemi v Python. Použití csv.reader s generátorem napájejícím 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()

Vzor generátoru udržuje využití paměti konstantní bez ohledu na velikost souboru. Pro mapování sloupců a zpracování identity viz Hromadné kopírovací operace.

Vložení jednoho řádku

Používejte jednotlivé inserty pro zápisy na úrovni aplikace, kdy zpracováváte jeden záznam najednou. Použití OUTPUT INSERTED pro získání generovaných klíčů:

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

Jednotlivé vložky jsou správnou volbou, když:

  • Na každou akci uživatele (odeslání formuláře, volání API) vložíte jeden řádek.
  • Každý řádek je potřeba ověřit nebo transformovat zvlášť před vložením.
  • Musíte okamžitě vložit ID nebo jiné vygenerované hodnoty.

Dávkové vložky

Použijte executemany() , když máte střední počet řádků a nepotřebujete průchodnost hromadné kopie:

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() odesílá každý řádek jako samostatný parametrizovaný příkaz. Když je propustnost důležitější než kontrola nad jednotlivými řádky, bulkcopy() je efektivnější, protože používá protokol TDS pro hromadné vkládání. Bod zvratu závisí na šířce řádků a latenci sítě, ale obvykle se pohybuje v řádu nižších stovek řádků.

Hromadné kopírování

Pokud je propustnost důležitější než řízení na řádek, použijte bulkcopy(). Používá protokol TDS bulk insert, který je výrazně efektivnější než vkládání řádek po řádku:

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

Tipy na výkon pro hromadné kopie

  • Používejte generátory pro velké datové sady, abyste udrželi konstantní využití paměti.
  • Nastavte batch_size pro určení, kolik řádků se odešle v jedné dávce TDS. Začni s 5 000 a uprav podle šířky řádku.
  • Používejte zámky na stoly pro speciální zatížení: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Deaktivujte indexy před načtením a poté je znovu sestavte. Tato sekvence zabraňuje režijním nákladům na údržbu indexů během načítání dat.

Pro mapování sloupců, identitní sloupce, zpracování NULL a paralelní načítání viz Hromadné kopírovací operace.

Upsert pomocí MERGE

MERGE je v jazyce Microsoft SQL příkaz pro podmíněné INSERT, UPDATE a DELETE v rámci jediné operace. Řeší vzor "vložit, pokud je nový, aktualizovat, pokud existuje", který Python vývojáři běžně potřebují.

Vložení nebo aktualizace jednoho řádku

Pro jeden řádek použijte MERGE s klauzulí USING , která definuje aliasy parametrů:

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

Objemový upsert s přípravnou tabulkou

Pro hromadné operace upsert nejprve nahrajte data do dočasné tabulky a potom je pomocí MERGE z ní aktualizujte. Použijte insert-or-update jako výchozí postup pro operace upsert u objektů DataFrame a pro dávkové aktualizace:

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

Tento příklad ukazuje výchozí vzor vkládání nebo aktualizace:

  • INSERT řádky ze zdroje, které v cílovém souboru neexistují (WHEN NOT MATCHED BY TARGET).
  • UPDATE řádky, které existují v obou (WHEN MATCHED).
  • Klauzule OUTPUT uvádí, jaké kroky byly provedeny v každém řádku, což je užitečné pro auditní stopy.

Caution

Přidávejte WHEN NOT MATCHED BY SOURCE THEN DELETE pouze tehdy, když jsou staging data autoritativním úplným snímkem cíle. Pokud batch obsahuje pouze změněné řádky, tato klauzule smaže řádky, které byly záměrně vynechány ze zdrojového feedu.

Pokud potřebujete úplnou synchronizaci, rozšiřte MERGE až poté, co potvrdíte, že zdroj je pro cílovou tabulku směrodatný:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

Ve sdílených prostředích používejte unikátní název globální dočasné tabulky pro každý běh nebo trvalou staging tabulku, abyste předešli kolizím mezi současnými úlohami.

Kdy místo toho použít samostatné příkazy UPDATE a INSERT

MERGE je výkonný, ale má krajní případy. Zvažte použití samostatných tvrzení, když:

  • Logiku nepotřebujete DELETE . Samostatné UPDATE, po kterém následuje INSERT WHERE NOT EXISTS, je čitelnější a snáze se ladí.
  • Tvrzení MERGE je natolik složité, že je těžké předvídat chování při zamykání. Samostatné příkazy vám umožňují explicitní kontrolu nad granularitou zámku.
  • Aktualizujete tabulku s vysokou souběžností, kde MERGE eskalace zámků může způsobit blokování.
# 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()

Načíst datové rámce

Extrahujte řádky z datového rámce Pandas nebo Polars pomocí bulkcopy():

pandas

Převeďte pandas DataFrame na n-tice a předejte do 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()

Polars

Převeďte Polars DataFrame na ntice pomocí metody .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()

Úplné způsoby načítání DataFramu najdete v integraci s pandas a integraci s Polars.

Parketové uspořádání

Použijte Parquet jako meziformát při migraci dat mezi systémy nebo když váš ETL pipeline již produkuje soubory 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()

U velkých souborů Parquet čtěte ve skupinách řádků, abyste zachovali konstantní využití paměti:

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

Validace načtených dat

Po načtení ověřte počet řádků a zkontrolujte data:

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

Pro produkční zátěže se nespoléhejte na transakci volajícího připojení k ochraně hovoru bulkcopy() . bulkcopy() otevře vlastní interní spojení a zkopírované řádky potvrdí nezávisle, takže je příkaz conn.rollback() na vašem hlavním spojení nemůže vrátit zpět. Dva přístupy vám dají atomicitu:

  • Nastavte use_internal_transaction=True tak, aby každá dávka byla zabalena do samostatné transakce. Dávka, která selže v průběhu zpracování, se vrátí do předchozího stavu, místo aby zůstala načtená jen z poloviny.
  • Chcete-li data před jejich přesunutím ověřit, hromadně je zkopírujte do pracovní tabulky, ověřte je a potom přesuňte řádky do cílové tabulky pomocí INSERT ... SELECT v rámci transakce v hlavním připojení. Protože to INSERT probíhá přes vaše připojení, conn.rollback() vrátí to zpět, pokud ověření selže.
# 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