Pilih pola pemuatan dan pemindahan data dengan mssql-python

Driver mssql-python menyediakan beberapa jalur untuk menulis data ke Microsoft SQL. Setiap jalur sesuai dengan beban kerja yang berbeda. Panduan ini membantu Anda memilih yang tepat berdasarkan volume data, format sumber, dan semantik pembaruan.

Tentukan berdasarkan beban kerja

Beban Kerja Jalur yang direkomendasikan Mengapa
Muat file CSV ke dalam tabel Muat data CSV dengan salinan massal bulkcopy() dengan generator menangani file dengan ukuran berapa pun tanpa memuatnya ke dalam memori.
Menyisipkan satu baris dari kode aplikasi Penyisipan satu baris Overhead yang rendah, penanganan kesalahan yang sederhana, kompatibel dengan OUTPUT untuk mengembalikan kunci yang dihasilkan.
Menyisipkan batch kecil hingga sedang dari kode aplikasi Penyisipan batch Mengurangi perjalanan pulang pergi dibandingkan dengan sisipan tunggal.
Muat ratusan baris atau lebih dari sumber mana pun Salinan massal Sisipan massal TDS adalah jalur paling efisien untuk volume besar.
Menyisipkan atau memperbarui baris berdasarkan kunci Upsert dengan MERGE MERGE menangani INSERT, UPDATE, dan DELETE dalam satu pernyataan.
Muat DataFrame ke dalam tabel Muat DataFrame Ekstrak baris dari pandas atau Polars dan teruskan ke bulkcopy().
Menyiapkan data melalui file Parquet Staging Parquet Berguna untuk ETL lintas sistem di mana format file perantara diperlukan.

Muat data CSV dengan salinan massal

Memasukkan data CSV adalah pertanyaan yang paling umum dalam pekerjaan basis data dengan Python. Gunakan csv.reader dengan generator yang menyuplai 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()

Pola generator menjaga penggunaan memori tetap konstan terlepas dari ukuran file. Untuk pemetaan kolom dan penanganan identitas, lihat Operasi penyalinan massal.

Sisipan baris tunggal

Gunakan sisipan tunggal untuk penulisan tingkat aplikasi di mana Anda memproses satu rekaman pada satu waktu. Gunakan OUTPUT INSERTED untuk mengambil kunci yang dihasilkan:

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

Sisipan tunggal adalah pilihan yang tepat ketika:

  • Anda menyisipkan satu baris per tindakan pengguna (pengiriman formulir, panggilan API).
  • Anda perlu memvalidasi atau mengubah setiap baris satu per satu sebelum menyisipkan.
  • Anda memerlukan ID yang dimasukkan atau nilai lain yang dihasilkan segera.

Penyisipan batch

Gunakan executemany() saat Anda memiliki jumlah baris sedang dan tidak memerlukan throughput salinan massal:

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() mengirimkan setiap baris sebagai pernyataan berparameter terpisah. Ketika laju pemrosesan lebih penting daripada kontrol tiap baris, bulkcopy() lebih efisien karena menggunakan protokol penyisipan massal TDS. Crossover bergantung pada lebar baris dan latensi jaringan, tetapi biasanya dalam ratusan baris yang rendah.

Penyalinan massal

Ketika throughput lebih penting daripada kontrol per baris, gunakan bulkcopy(). Ini menggunakan protokol sisipan massal TDS, yang secara signifikan lebih efisien daripada sisipan baris demi baris:

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

Tips performa untuk salinan massal

  • Gunakan generator untuk himpunan data besar agar penggunaan memori tetap konstan.
  • Mengatur batch_size untuk mengontrol berapa banyak baris yang dikirim per batch TDS. Mulailah dengan 5.000 dan sesuaikan berdasarkan lebar baris.
  • Gunakan kunci meja untuk beban eksklusif: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Nonaktifkan indeks sebelum memuat, lalu bangun kembali setelahnya. Urutan ini menghindari overhead pemeliharaan indeks selama proses pemuatan.

Untuk pemetaan kolom, kolom identitas, penanganan NULL, dan pemuatan paralel, lihat Operasi penyalinan massal.

Upsert dengan MERGE

MERGEadalah pernyataan Microsoft SQL untuk bersyarat INSERT, UPDATE, dan DELETE dalam satu operasi. Ini menangani pola "sisipkan jika baru, perbarui jika ada" yang biasanya dibutuhkan pengembang Python.

Upsert baris tunggal

Untuk satu baris, gunakan MERGE dengan USING klausa yang mendefinisikan alias parameter:

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

Upsert massal dengan tabel pementasan

Untuk upsert massal, muat data ke tabel sementara terlebih dahulu, lalu gunakan MERGE untuk melakukan pembaruan dari tabel tersebut. Gunakan insert-or-update sebagai pola default untuk upsert DataFrame dan pembaruan batch:

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

Contoh ini menunjukkan pola sisipkan atau pembaruan default:

  • INSERT baris dari sumber yang tidak ada di target (WHEN NOT MATCHED BY TARGET).
  • UPDATE baris yang ada di keduanya (WHEN MATCHED).
  • Klausul OUTPUT melaporkan tindakan apa yang diambil pada setiap baris, yang berguna untuk jejak audit.

Caution

Tambahkan WHEN NOT MATCHED BY SOURCE THEN DELETE hanya jika data penahapan adalah rekam jepret penuh target yang otoritatif. Jika batch hanya berisi baris yang diubah, klausa tersebut menghapus baris yang sengaja dihilangkan dari umpan sumber.

Jika Anda memerlukan rekonsiliasi penuh, perluas MERGE hanya setelah Anda mengonfirmasi bahwa sumber tersebut merupakan sumber yang otoritatif untuk tabel target:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

Di lingkungan bersama, gunakan nama tabel sementara global yang unik untuk setiap eksekusi atau tabel staging permanen guna menghindari konflik antarpekerjaan yang berjalan secara bersamaan.

Kapan harus menggunakan pernyataan terpisah UPDATE dan INSERT sebagai gantinya

MERGE Kuat tetapi memiliki casing tepi. Pertimbangkan untuk menggunakan pernyataan terpisah ketika:

  • Anda tidak membutuhkan DELETE logika. Penggunaan UPDATE secara terpisah yang diikuti oleh INSERT WHERE NOT EXISTS lebih mudah dibaca serta lebih mudah di-debug.
  • Pernyataan MERGE ini cukup kompleks sehingga perilaku pengunciannya sulit diprediksi. Pernyataan terpisah memberi Anda kontrol eksplisit atas granularitas kunci.
  • Anda memperbarui tabel konkurensi tinggi di mana MERGE eskalasi kunci dapat menyebabkan pemblokiran.
# 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()

Muat DataFrame

Ekstrak baris dari pandas atau Polars DataFrame dan muat dengan menggunakan bulkcopy():

pandas

Konversi DataFrame pandas menjadi tuple dan teruskan ke 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

Mengonversi DataFrame Polars menjadi tuple menggunakan metode .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()

Untuk pola pemuatan DataFrame lengkap, lihat integrasi pandas dan integrasi Polars.

Pementasan parket

Gunakan Parquet sebagai format perantara saat memigrasikan data antar sistem atau saat alur ETL Anda sudah menghasilkan file 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()

Untuk file Parquet besar, baca dalam grup baris agar penggunaan memori tetap konstan:

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

Memvalidasi data yang dimuat

Setelah dimuat, verifikasi jumlah baris dan lakukan pemeriksaan sampel pada 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}")

Untuk beban kerja produksi, jangan mengandalkan transaksi koneksi pemanggil untuk melindungi panggilan bulkcopy(). bulkcopy() membuka koneksi internalnya sendiri dan melakukan commit atas baris yang disalin secara independen, sehingga conn.rollback() pada koneksi utama Anda tidak dapat membatalkan perubahan tersebut. Dua pendekatan memberi Anda atomisitas:

  • Atur use_internal_transaction=True untuk membungkus setiap batch dalam transaksinya masing-masing. Batch yang gagal setengah jalan menggulirkan kembali batch tersebut alih-alih membiarkannya setengah dimuat.
  • Untuk memvalidasi data sebelum menaikkannya ke tahap produksi, salin data secara massal ke tabel staging, validasi data tersebut, lalu pindahkan baris-barisnya ke tabel target dengan menggunakan INSERT ... SELECT di dalam transaksi pada koneksi utama Anda. Karena itu INSERT berjalan pada koneksi Anda, conn.rollback() urungkan jika validasi gagal.
# 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