Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Perpustakaan pandas adalah alat analisis data utama Python. Dengan menggabungkan panda dengan driver mssql-python, Anda dapat:
- Muat hasil kueri SQL langsung ke DataFrames.
- Tulis DataFrames kembali ke Microsoft SQL secara efisien.
- Lakukan operasi ETL.
- Buat alur data.
Contoh-contoh dalam artikel ini mengkueri tabel Production.Product dan tabel lainnya di database sampel AdventureWorks. Contoh yang menulis data menggunakan tabel sementara untuk menghindari memodifikasi data sampel.
Tabel lain yang direferensikan dalam contoh analisis (Sales.SalesOrderHeader, Sales.SalesOrderDetail, ) Production.ProductSubcategoryadalah bagian dari AdventureWorks. Gunakan tabel Anda sendiri sebagai pengganti saat menyesuaikan pola ini.
Membaca data ke dalam DataFrames
Driver mssql-python mengembalikan baris sebagai objek Python, yang Anda konversi menjadi pandas DataFrames dengan membaca nama kolom dari cursor.description dan nilai baris dari fetchall(). Fungsi pembantu di bagian ini membungkus konversi itu menjadi pola yang dapat digunakan kembali.
Kueri dasar ke DataFrame
Fungsi ini mengeksekusi kueri berparameter dan membangun DataFrame dari kumpulan hasil lengkap. Ini berfungsi dengan baik untuk kumpulan hasil yang dapat ditampung dengan nyaman di memori.
import pandas as pd
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
def query_to_dataframe(cursor, query: str, params: dict = None) -> pd.DataFrame:
"""Execute query and return results as DataFrame."""
cursor.execute(query, params or {})
# cursor.description is a list of tuples, one per column.
# Each tuple's first element is the column name.
columns = [col[0] for col in cursor.description]
# Fetch all rows
rows = cursor.fetchall()
# Convert to DataFrame
data = [tuple(row) for row in rows]
return pd.DataFrame(data, columns=columns)
# Usage: %(cat)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_dataframe(cursor, "SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 5})
print(df.head())
Note
Jika string koneksi Anda menggunakan Authentication=ActiveDirectoryDefault, driver menggunakan DefaultAzureCredential, yang mencoba beberapa penyedia kredensial secara berurutan. Koneksi pertama bisa lambat karena SDK menelusuri rantai hingga menemukan penyedia yang berfungsi. Dalam lingkungan produksi, jika Anda mengetahui jenis kredensial yang digunakan oleh lingkungan Anda, tentukan secara langsung (misalnya, ActiveDirectoryMSI untuk identitas terkelola) agar terhindar dari penelusuran berantai. Untuk informasi selengkapnya, lihat Autentikasi Microsoft Entra.
Mengalirkan himpunan data besar
Untuk tabel dengan jutaan baris, memuat semuanya sekaligus dapat menghabiskan memori. Pendekatan berbasis potongan mengambil baris secara bertahap dalam batch dengan fetchmany() dan menggabungkan hasilnya, sehingga penggunaan memori puncak tetap sebanding dengan chunksize, bukan dengan seluruh kumpulan hasil.
def query_to_dataframe_chunked(cursor, query: str, params: dict = None,
chunksize: int = 10000) -> pd.DataFrame:
"""Load large query results in chunks for memory efficiency."""
cursor.execute(query, params or {})
columns = [col[0] for col in cursor.description]
chunks = []
while True:
rows = cursor.fetchmany(chunksize)
if not rows:
break
data = [tuple(row) for row in rows]
chunks.append(pd.DataFrame(data, columns=columns))
return pd.concat(chunks, ignore_index=True) if chunks else pd.DataFrame(columns=columns)
# Usage for large tables
df = query_to_dataframe_chunked(cursor, "SELECT * FROM Production.TransactionHistory", chunksize=50000)
Generator untuk himpunan data besar
Saat Anda perlu memproses data secara bertahap tanpa menyimpan seluruh hasil dalam memori, gunakan generator. Masing-masing yield menghasilkan satu potongan DataFrame yang dapat Anda proses dan buang sebelum mengambil yang berikutnya.
def query_to_dataframe_generator(cursor, query: str, params: dict = None,
chunksize: int = 10000):
"""Yield DataFrame chunks for processing without loading all data."""
cursor.execute(query, params or {})
columns = [col[0] for col in cursor.description]
while True:
rows = cursor.fetchmany(chunksize)
if not rows:
break
data = [tuple(row) for row in rows]
yield pd.DataFrame(data, columns=columns)
# Process chunks without loading entire dataset
huge_query = """
SELECT * FROM Production.TransactionHistory
UNION ALL SELECT * FROM Production.TransactionHistory
UNION ALL SELECT * FROM Production.TransactionHistory
"""
for chunk_df in query_to_dataframe_generator(cursor, huge_query):
# Process each chunk, then discard it before the next fetch
print(f"Processing chunk of {len(chunk_df)} rows")
total_cost = chunk_df["ActualCost"].sum()
print(f"Chunk total cost: {total_cost}")
Menulis DataFrames ke Microsoft SQL
Pengidentifikasi kutipan untuk mencegah injeksi SQL
Nama tabel dan kolom tidak dapat diteruskan sebagai parameter kueri di SQL. Saat Anda membuat pernyataan SQL dengan pengidentifikasi dinamis, bungkus setiap nama dalam tanda kurung siku dan keluarkan karakter yang disematkan ] untuk mencegah injeksi SQL.
def quote_id(identifier: str) -> str:
"""Quote a Microsoft SQL identifier to prevent SQL injection.
Wraps the name in square brackets and escapes any embedded ] characters.
Raises ValueError if the identifier is empty or contains null bytes.
"""
if not identifier or "\x00" in identifier:
raise ValueError(f"Invalid identifier: {identifier!r}")
escaped = identifier.replace("]", "]]")
return f"[{escaped}]"
Fungsi pembantu di bagian ini digunakan quote_id() untuk semua nama tabel dan kolom di SQL yang dihasilkan.
Sisipkan baris DataFrame
Pendekatan paling sederhana melakukan iterasi pada setiap baris DataFrame dan menerbitkan satu INSERT per baris. Pendekatan sederhana berfungsi untuk DataFrames kecil tetapi lambat untuk volume besar karena setiap baris memerlukan perjalanan pulang pergi terpisah ke server.
def dataframe_to_sql(cursor, conn, df: pd.DataFrame, table: str,
if_exists: str = "append") -> int:
"""Write DataFrame to Microsoft SQL table."""
if if_exists == "replace":
cursor.execute(f"TRUNCATE TABLE {quote_id(table)}")
columns = df.columns.tolist()
placeholders = ", ".join([f"%({col})s" for col in columns])
col_list = ", ".join([quote_id(col) for col in columns])
query = f"INSERT INTO {quote_id(table)} ({col_list}) VALUES ({placeholders})"
rows_inserted = 0
for _, row in df.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
cursor.execute(query, params)
rows_inserted += 1
conn.commit()
return rows_inserted
# Usage
cursor.execute("""
CREATE TABLE #Products (
Name NVARCHAR(100),
ListPrice DECIMAL(10,2),
ProductSubcategoryID INT
)
""")
df = pd.DataFrame({
"Name": ["Product A", "Product B"],
"ListPrice": [29.99, 49.99],
"ProductSubcategoryID": [1, 2]
})
rows = dataframe_to_sql(cursor, conn, df, "#Products")
print(f"Inserted {rows} rows")
Penyisipan massal dengan BCP (direkomendasikan untuk DataFrame berukuran besar)
Untuk DataFrame berukuran besar, gunakan metode bulkcopy() driver, yang mengirimkan baris secara batch melalui protokol TDS (Tabular Data Stream), yaitu protokol wire native yang digunakan oleh Microsoft SQL. Pendekatan ini lebih cepat daripada sisipan baris demi baris karena meminimalkan perjalanan pulang pergi.
def dataframe_to_sql_bulk(conn, df: pd.DataFrame, table: str) -> int:
"""Bulk insert DataFrame using BCP for better performance."""
# Convert DataFrame to list of tuples, handling NaN
rows = []
for _, row in df.iterrows():
row_data = tuple(None if pd.isna(v) else v for v in row)
rows.append(row_data)
cursor = conn.cursor()
result = cursor.bulkcopy(table, rows)
conn.commit()
return result["rows_copied"]
# Usage
cursor.execute("CREATE TABLE ##PandasProducts (Name NVARCHAR(50), ListPrice DECIMAL(10,2), ProductSubcategoryID INT)")
conn.commit()
df = pd.DataFrame({
"Name": ["Product A", "Product B", "Product C"],
"ListPrice": [29.99, 49.99, 19.99],
"ProductSubcategoryID": [1, 2, 1]
})
rows = dataframe_to_sql_bulk(conn, df, "##PandasProducts")
Perbarui baris yang ada dari DataFrame
Untuk memperbarui baris yang sudah ada di dalam tabel, lakukan iterasi pada DataFrame dan jalankan pernyataan UPDATE berparameter.
key_column mengidentifikasi baris yang akan diperbarui.
def update_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
key_column: str) -> int:
"""Update existing rows based on key column."""
columns = [col for col in df.columns if col != key_column]
set_clause = ", ".join([f"{quote_id(col)} = %({col})s" for col in columns])
query = f"UPDATE {quote_id(table)} SET {set_clause} WHERE {quote_id(key_column)} = %({key_column})s"
rows_updated = 0
for _, row in df.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
cursor.execute(query, params)
rows_updated += cursor.rowcount
conn.commit()
return rows_updated
# Usage
cursor.execute("""
CREATE TABLE #ProductPrices (
ProductID INT PRIMARY KEY,
ListPrice DECIMAL(10,2)
);
INSERT INTO #ProductPrices VALUES (1, 29.99), (2, 49.99), (3, 19.99);
""")
conn.commit()
df_updates = pd.DataFrame({
"ProductID": [1, 2, 3],
"ListPrice": [31.99, 52.99, 21.99]
})
updated = update_from_dataframe(cursor, conn, df_updates, "#ProductPrices", "ProductID")
Pola upsert (penggabungan)
Ketika beberapa baris mungkin baru dan yang lain mungkin sudah ada, gunakan pernyataan SQL MERGE untuk menyisipkan atau memperbarui dalam satu operasi.
MERGE membandingkan setiap baris masuk dengan tabel target menggunakan kolom kunci. Jika ditemukan kecocokan, data akan diperbarui; jika tidak, data akan disisipkan.
MERGE menghindari pemeriksaan keberadaan secara terpisah.
def upsert_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
key_columns: list[str]) -> int:
"""Insert or update rows based on key columns. Returns total rows affected."""
all_columns = df.columns.tolist()
value_columns = [c for c in all_columns if c not in key_columns]
total_affected = 0
for _, row in df.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
# Build MERGE statement with quoted identifiers
key_match = " AND ".join([f"t.{quote_id(k)} = s.{quote_id(k)}" for k in key_columns])
update_set = ", ".join([f"{quote_id(c)} = s.{quote_id(c)}" for c in value_columns])
all_cols = ", ".join([quote_id(c) for c in all_columns])
all_vals = ", ".join([f"%({c})s" for c in all_columns])
cursor.execute(f"""
MERGE {quote_id(table)} AS t
USING (SELECT {', '.join([f'%({c})s AS {quote_id(c)}' for c in all_columns])}) AS s
ON {key_match}
WHEN MATCHED THEN UPDATE SET {update_set}
WHEN NOT MATCHED THEN INSERT ({all_cols}) VALUES ({all_vals});
""", params)
total_affected += cursor.rowcount
conn.commit()
return total_affected
Pola analisis data
Contoh berikut menunjukkan tugas analisis umum yang menggabungkan kueri Microsoft SQL dengan transformasi pandas.
Menggabungkan kueri ke DataFrame
def get_sales_summary(cursor) -> pd.DataFrame:
"""Get sales summary by category."""
return query_to_dataframe(cursor, """
SELECT
pc.Name AS CategoryName,
COUNT(*) AS ProductCount,
AVG(p.ListPrice) AS AvgPrice,
MIN(p.ListPrice) AS MinPrice,
MAX(p.ListPrice) AS MaxPrice
FROM Production.Product p
JOIN Production.ProductSubcategory pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
GROUP BY pc.Name
ORDER BY ProductCount DESC
""")
df = get_sales_summary(cursor)
print(df.to_string())
Data rangkaian waktu
Gunakan pengindeksan tanggal dan resampling pandas untuk mengolah data deret waktu dari Microsoft SQL. Untuk mengaktifkan operasi seperti rata-rata bergulir dan pengambilan sampel ulang, atur kolom tanggal sebagai indeks DataFrame.
def get_daily_sales(cursor, start_date: str, end_date: str) -> pd.DataFrame:
"""Get daily sales time series."""
df = query_to_dataframe(cursor, """
SELECT
CAST(OrderDate AS DATE) AS Date,
COUNT(*) AS OrderCount,
SUM(TotalDue) AS Revenue
FROM Sales.SalesOrderHeader
WHERE OrderDate BETWEEN %(start)s AND %(end)s
GROUP BY CAST(OrderDate AS DATE)
ORDER BY Date
""", {"start": start_date, "end": end_date})
# Set date as index for time series operations
df["Date"] = pd.to_datetime(df["Date"])
df.set_index("Date", inplace=True)
return df
# Usage
sales_df = get_daily_sales(cursor, "2024-01-01", "2024-12-31")
# Resample to weekly
weekly = sales_df.resample("W").sum()
# Calculate rolling average
sales_df["RollingAvg"] = sales_df["Revenue"].rolling(window=7).mean()
Tabel pivot dari data SQL
Tabel pivot membentuk ulang data dari baris menjadi format matriks. Untuk mengaturnya ulang berdasarkan dimensi seperti tahun, bulan, dan kategori, tarik data mentah dari Microsoft SQL, lalu gunakan pivot_table().
def get_sales_pivot(cursor) -> pd.DataFrame:
"""Get sales data and create pivot table."""
df = query_to_dataframe(cursor, """
SELECT
YEAR(soh.OrderDate) AS Year,
MONTH(soh.OrderDate) AS Month,
pc.Name AS CategoryName,
SUM(sod.OrderQty * sod.UnitPrice) AS Revenue
FROM Sales.SalesOrderHeader soh
JOIN Sales.SalesOrderDetail sod ON soh.SalesOrderID = sod.SalesOrderID
JOIN Production.Product p ON sod.ProductID = p.ProductID
JOIN Production.ProductSubcategory pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
GROUP BY YEAR(soh.OrderDate), MONTH(soh.OrderDate), pc.Name
""")
# Create pivot table
pivot = df.pivot_table(
values="Revenue",
index=["Year", "Month"],
columns="CategoryName",
aggfunc="sum",
fill_value=0
)
return pivot
pivot_df = get_sales_pivot(cursor)
print(pivot_df)
Pola ETL
Untuk membangun alur ekstrak, transformasi, dan muat, gabungkan kueri Microsoft SQL dengan transformasi pandas untuk membangun alur ekstrak, transformasi, dan muat. Pengemudi menangani ekstraksi dan pemuatan sementara panda menangani langkah transformasi.
Mengekstrak, mengubah, memuat
Contoh ini mengekstrak data pelanggan aktif, menerapkan aturan bisnis untuk mengelompokkan pelanggan, dan memuat hasilnya ke dalam tabel tujuan.
def etl_pipeline(source_cursor, dest_cursor, dest_conn):
"""Simple ETL pipeline with pandas."""
# Extract
df = query_to_dataframe(source_cursor, """
SELECT
c.CustomerID,
COUNT(soh.SalesOrderID) AS OrderCount,
SUM(soh.TotalDue) AS TotalSpent
FROM Sales.Customer c
JOIN Sales.SalesOrderHeader soh ON c.CustomerID = soh.CustomerID
WHERE soh.OrderDate > DATEADD(YEAR, -1, GETDATE())
GROUP BY c.CustomerID
""")
# Transform
df["CustomerSegment"] = pd.cut(
df["TotalSpent"],
bins=[0, 100, 500, 1000, float("inf")],
labels=["Bronze", "Silver", "Gold", "Platinum"]
)
df["AvgOrderValue"] = df["TotalSpent"] / df["OrderCount"].replace(0, 1)
df["IsHighValue"] = df["TotalSpent"] > 500
# Load
dataframe_to_sql_bulk(dest_conn, df[["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"]],
"#CustomerAnalytics")
return len(df)
Pola beban tambahan
Untuk alur data yang sedang berlangsung, muat hanya rekaman yang berubah sejak eksekusi terakhir. Pendekatan ini mengkueri tabel tujuan untuk stempel waktu maksimum, lalu hanya mengambil rekaman yang lebih baru dari sumber.
def incremental_load(cursor, conn, source_table: str, dest_table: str,
timestamp_col: str) -> int:
"""Load only new/changed records based on timestamp."""
# Get last loaded timestamp
cursor.execute(f"SELECT MAX({quote_id(timestamp_col)}) FROM {quote_id(dest_table)}")
last_loaded = cursor.fetchval()
# Build query for new records
if last_loaded:
df = query_to_dataframe(cursor, f"""
SELECT * FROM {quote_id(source_table)}
WHERE {quote_id(timestamp_col)} > %(last)s
""", {"last": last_loaded})
else:
df = query_to_dataframe(cursor, f"SELECT * FROM {quote_id(source_table)}")
if df.empty:
return 0
# Load new records
return dataframe_to_sql_bulk(conn, df, dest_table)
Tips kinerja
Menggunakan jenis data yang sesuai
Secara bawaan, Pandas menggunakan tipe 64-bit untuk data numerik, sehingga memboroskan memori meskipun tipe yang lebih kecil sudah memadai. Menurunkan bilangan bulat dan float, dan mengonversi kolom string kardinalitas rendah menjadi kategoris, dapat secara signifikan mengurangi penggunaan memori.
def optimize_dataframe_types(df: pd.DataFrame) -> pd.DataFrame:
"""Optimize DataFrame memory usage."""
for col in df.columns:
col_type = df[col].dtype
if col_type == "int64":
# Downcast integers
df[col] = pd.to_numeric(df[col], downcast="integer")
elif col_type == "float64":
# Downcast floats
df[col] = pd.to_numeric(df[col], downcast="float")
elif col_type == "object":
# Convert to category if low cardinality
num_unique = df[col].nunique()
if num_unique / len(df) < 0.5:
df[col] = df[col].astype("category")
return df
Gunakan SQL untuk mengangkat beban berat
Microsoft SQL lebih cepat untuk agregasi, pemfilteran, dan gabungan daripada menarik semua data mentah Anda melalui kabel dan pemrosesan secara lokal di Python. Biarkan Microsoft SQL melakukan pekerjaan berat jika memungkinkan, pindahkan hanya data yang Anda butuhkan di seluruh jaringan, dan gunakan panda untuk analisis dan transformasi yang lebih nyaman di Python.
# Avoid: pulling all rows over the wire to aggregate locally in pandas
df_all = query_to_dataframe(cursor, "SELECT * FROM Production.Product") # transfers entire table
summary = df_all.groupby("Color").agg({"ListPrice": "sum"}) # aggregation that SQL can do faster
# Better: push the aggregation into SQL and transfer only the summary
df = query_to_dataframe(cursor, """
SELECT Color, SUM(ListPrice) AS TotalPrice
FROM Production.Product
WHERE Color IS NOT NULL
GROUP BY Color
""")
Penulisan secara batch
Untuk DataFrame besar yang terlalu besar untuk satu sisipan massal, pisahkan pekerjaan menjadi beberapa batch dan lacak kemajuan.
def batch_insert(cursor, conn, df: pd.DataFrame, table: str, batch_size: int = 1000):
"""Insert in batches with progress tracking."""
total = len(df)
for i in range(0, total, batch_size):
batch = df.iloc[i:i + batch_size]
dataframe_to_sql(cursor, conn, batch, table)
print(f"Inserted {min(i + batch_size, total)}/{total}")