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 bulkcopy_arrow()olvassa a DataFrame Arrow adatait anélkül, hogy minden értékhez Python objektumot építene.
Apache Arrow adatait töltsd be egy táblázatba Arrow-adatok betöltése bulkcopy_arrow()közvetlenül olvassa az Arrow memóriáját, anélkül, hogy Python tuple-okat építene.
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. A Parquet eleve Arrow-formátumú adat, így soronkénti átalakítás nélkül betöltődik.

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 beillesztési protokollt használja, amely sorokat streamel, ahelyett, hogy soronként egy utasítást küldene:

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.
  • Használja a(z) bulkcopy_arrow() elemet, ha a forrás oszlopalapú, például egy DataFrame vagy egy Parquet-fájl. Kihagyja a Pythonbeli sor-tuple-ökké alakítást.
  • Á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

A DataFrame oszlopalapú, ezért bulkcopy_arrow() töltsd be, ahelyett, hogy a bulkcopy() számára sor-tuple-ökre lapítanád.

Egyeztesd össze az Arrow típusokat a céloszlopokkal a betöltés előtt. pyarrow egy float64 numerikus oszlopra következtet, amelyet a meghajtó nem tud pénzre, tizedesre vagy numerikusra leképezni.

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

Használd inkább a(z) Table.cast() elemet, mint hogy a sémát átadd a(z) Table.from_pandas() elemnek, amely nem tud egy lebegőpontos oszlopot közvetlenül decimal128 típussá alakítani.

Polarok

A Polars megvalósítja az Arrow C adatinterfészt, így át tudod adni magát a DataFrame-et. Először az oszlopokat öntjük ugyanazért az okból:

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

A típusokat a fájl beolvasásakor is beállíthatod a pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}) használatával.

A DataFrame átadása közvetlenül a meghajtónak adja át a puffereket másolat nélkül. df.to_arrow() szintén működik, de a Polars a konverzió során újra kódolja a string oszlopokat, ami az összes string adatot lemásolja.

bulkcopy_arrow() elfogad egy pyarrow.Table, egy RecordBatch, egy RecordBatchReader elemet, vagy bármely objektumot, amely a __arrow_c_stream__ vagy __arrow_c_array__ révén megvalósítja az Arrow C adatinterfészt. Ha ezek közül bármelyiket átadjuk a bulkcopy()-nak, az TypeError kivételt vált ki.

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

Arrow-adatok betöltése

Ha a forrás már Apache Arrow formátumban van, cursor.bulkcopy_arrow() betölti anélkül, hogy először a Python tuple-okat kellene építeni.

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

A módszer egy pyarrow.RecordBatch vagy egy pyarrow.RecordBatchReader elemet is fogad, így egy eredményhalmaz a(z) cursor.arrow_reader()-ből közvetlenül egy másik táblába továbbítható.

Minden Arrow oszloptípusnak kompatibilisnek kell lennie a célállomás SQL oszloptípusával, és az író nem konvertál típuscsaládok között. További információért lásd: Apache Arrow 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. A Parquet-fájl Arrow-adatként olvasható be, ezért add át közvetlenül a bulkcopy_arrow() számára:

import pyarrow.parquet as pq

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

Nagy Parquet fájlok esetén iterálj sorcsoportokat, hogy állandó memóriahasználat maradjon. Minden köteg egy RecordBatch, amelyet a bulkcopy_arrow() közvetlenül elfogad:

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

Ha a teljes fájlt egyetlen hívással szeretné streamelni, csomagolja a kötegeket a RecordBatchReader elembe:

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

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