Válassz adatbetöltési és mozgási mintát mssql-python segítségével

Az illezőprogram mssql-python több útvonalat biztosít az adatok Microsoft SQL-be történő írásához. Minden útvonal más-más munkaterheléshez illeszkedik. Ez az útmutató segít kiválasztani a megfelelőt az adatmennyiség, a forrásformátum és a friss szemantika alapján.

Döntés a munkaterhelés szerint

Munkaterhelés Ajánlott elérési út Miért
CSV fájlok betöltése egy táblába CSV-adatok betöltése tömeges másolással bulkcopy() a generátor bármilyen méretű fájlokat kezel anélkül, hogy betöltené őket a memóriába.
Egyetlen sor beszúrása alkalmazáskódból Egysoros beszúrások Alacsony túlterhelés, egyszerű hibakezelés, működik a OUTPUT segítségével a generált kulcsok visszaadásához.
Kis-közepes adag beépítése alkalmazáskódból Kötegelt beszúrások Csökkenti a oda-vissza utazásokat az egyes betétekhez képest.
Tölts be több száz vagy akár ennél is több sort bármilyen forrásból tömeges átmásolás A TDS tömeges beillesztés a leghatékonyabb út nagy térfogatok esetén.
Sorok beépítése vagy frissítése kulcs alapján Upsert ezzel: MERGE MERGE egy utasításban kezeli a(z) INSERT, UPDATE és DELETE elemeket.
Adatkeret betöltése egy táblába Adatkeretek betöltése Nyerje ki a sorokat a pandasból vagy a Polarsból, és adja át őket a(z) bulkcopy() számára.
Adatok átmeneti tárolása Parquet-fájlok használatával Parquet átmeneti tárolás Hasznos rendszerközi ETL-hez, ahol köztes fájlformátumra van szükség.

CSV-adatok betöltése tömeges másolással

A CSV-adatok betöltése a leggyakoribb adatbetöltéssel kapcsolatos kérdés a Pythonos adatbázis-kezelésben. Használja a(z) csv.reader elemet, ha egy generátor táplálja a(z) bulkcopy() elemet:

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

A generátor minta állandóan tartja a memóriahasználatot a fájlmérettől függetlenül. Az oszlopleképezés és az identitás kezeléséről lásd: Tömeges másolási műveletek.

Egysoros beszúrások

Használj egyetlen beszedést alkalmazásszintű írásokhoz, ahol egyszerre csak egy rekordot dolgozol fel. A generált kulcsok visszakeresésére használják 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()

Az egyszemélyes betétek a megfelelő választás, ha:

  • Minden felhasználói művelethez (űrlap elküldése, API-hívás) egy sort szúrsz be.
  • Minden sort külön-külön kell validálni vagy átalakítani a beillesztés előtt.
  • Azonnal szükséged van a beillesztett azonosítóra vagy más generált értékekre.

Kötegelt beszúrások

Használd executemany(), ha közepes számú sorral dolgozol, és nincs szükséged a tömeges másolás áteresztőképességére:

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() minden sort külön paraméterezett állításként küld. Amikor az áteresztés fontosabb, mint soronkénti vezérlés, bulkcopy() hatékonyabb, mert a TDS tömeges beillesztési protokollt használja. Az átváltási pont a sorok szélességétől és a hálózati késleltetéstől függ, de jellemzően néhány száz sornál van.

Tömeges másolás

Ha az átviteli sebesség fontosabb, mint a soronkénti vezérlés, használd a bulkcopy() elemet. A TDS tömeges beszúrási protokollt használja, amely jelentősen hatékonyabb, mint a soronkénti beszúrások:

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

Teljesítménytippek tömeges másolathoz

  • Használj generátorokat nagy adathalmazokhoz, hogy állandóan tartsd a memóriahasználatot.
  • Állítsa be a(z) batch_size értékét annak szabályozásához, hogy hány sort küldjön el egy TDS-kötegben. Kezdj 5 000-vel, és állítsd be a sorszélesség alapján.
  • Használj asztalzárakat exkluzív terhelésekhez: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Kapcsold ki az indexeket a betöltés előtt, majd utána építsd újra. Ez a sorozat elkerüli az indexfenntartási túlterhelést a terhelés alatt.

Oszlopleképezések, identitásoszlopok, NULL kezelés és párhuzamos betöltés esetén lásd: Tömeges másolási műveletek.

Upsert ezzel: MERGE

MERGE a Microsoft SQL feltételes INSERT, UPDATE és DELETE utasítása egyetlen műveletben. Kezeli azt a "beépíts, ha új, frissítsd, ha létezik" mintát, amire a Python fejlesztőknek gyakran szükségük van.

Egysoros felsúr

Egyetlen sor esetén használja a MERGE elemet egy USING záradékkal, amely paraméterálneveket definiál:

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

Tömeges upsert átmeneti táblával

Tömeges upsert műveletekhez először töltsd be az adatokat egy ideiglenes táblába, majd a MERGE használatával frissíts onnan. Használd az insert-or-update alapértelmezett mintát DataFrame upsert-ekhez és batch frissítésekhez:

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

Ez a példa bemutatja az alapértelmezett beillesztés vagy frissítés mintát:

  • INSERT olyan sorok a forrásból, amelyek nem léteznek a célpontban (WHEN NOT MATCHED BY TARGET).
  • UPDATE mindkettőben megtalálható sorok (WHEN MATCHED).
  • A OUTPUT záradék bemutatja, hogy milyen lépéseket tettek az egyes sorokban, ami hasznos az audit nyomok során.

Caution

A WHEN NOT MATCHED BY SOURCE THEN DELETE elemet csak akkor add hozzá, ha a staging adatok a cél teljes, hiteles pillanatképét jelentik. Ha a batch csak megváltoztatott sorokat tartalmaz, az a záradék törli azokat a sorokat, amelyeket szándékosan kihagytak a forrásforrásból.

Ha teljes egyeztetésre van szükséged, a MERGE elemet csak azután terjeszd ki, hogy megerősítetted: a forrás hiteles a céltábla szempontjából:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

Megosztott környezetekben futtatásonként egyedi globális ideiglenestábla-nevet, illetve az egyidejűleg futó feladatok közötti ütközések elkerülése érdekében állandó átmeneti táblát használjon.

Mikor kell különálló UPDATE és INSERT állításokat használni helyette

MERGE erőteljes, de vannak szélsőséges esetei. Érdemes külön állításokat használni, ha:

  • Nincs szükséged DELETE logikára. A különálló UPDATE, amelyet INSERT WHERE NOT EXISTS követ, olvashatóbb, és könnyebb hibakeresni.
  • Az állítás MERGE elég összetett ahhoz, hogy a zárolási viselkedést nehéz előre látni. Különálló állítások explicit irányítást adnak a zár granularitása felett.
  • Egy magas párhuzamú táblázatot frissítesz, ahol MERGE a zár ekalációja blokkolást okozhat.
# 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()

Adatkeretek betöltése

Sorokat húzz ki egy pandas vagy Polars DataFrame-ből, és töltsd be őket a következők segítségével bulkcopy():

pandas

A pandas DataFrame átalakítása tuple-ökké és átadása a következőnek: 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()

Polarok

A Polars DataFrame átalakítása tuple-okká a következő .rows() módszerrel:

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

A teljes DataFrame betöltési mintákért lásd: pandas integráció és Polars integráció.

Parquet átmeneti tárolás

Használd a Parquet-et köztes formátumként, amikor adatokat migrálsz rendszerek között, vagy amikor az ETL csővezeték már Parquet fájlokat generál:

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

Nagy Parquet fájlok esetén sorcsoportokban olvassuk, hogy állandó memóriahasználat maradjon:

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

Validáld a betöltött adatokat

A betöltés után ellenőrizd a sorok számát, és szúrópróbaszerűen ellenőrizd az adatokat:

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

Éles terhelésnél ne támaszkodj a hívást indító kapcsolat tranzakciójára a bulkcopy() hívás védelméhez. bulkcopy() megnyitja a saját belső kapcsolatát, és a másolt sorokat önállóan véglegesíti, így a fő kapcsolatodon lévő conn.rollback() nem tudja ezeket visszavonni. Két megközelítés ad atomitást:

  • Úgy állítsuk use_internal_transaction=True be, hogy minden tételt saját tranzakcióba csomagolj. Az a köteg, amelynek a feldolgozása menet közben meghiúsul, visszagörgeti az adott köteget ahelyett, hogy félig betöltött állapotban hagyná.
  • Az adatok célhelyre helyezése előtti érvényesítéséhez másold őket tömegesen egy átmeneti táblába, érvényesítsd őket, majd a fő kapcsolaton indított tranzakción belül egy INSERT ... SELECT használatával helyezd át a sorokat a céltáblába. Mivel ez INSERT az Ön kapcsolatán keresztül fut, conn.rollback() visszavonja azt, ha az ellenőrzés sikertelen.
# 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