Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Sterownik mssql-python zawiera funkcję kopiowania masowego, która efektywnie wstawia duże ilości danych do SQL Server, Azure SQL Database, Azure SQL Managed Instance oraz bazy danych SQL w Microsoft Fabric.
Metoda cursor.bulkcopy() oferuje wydajny sposób wczytywania dużych zbiorów danych:
- Minimalizuje to podróże sieciowe w obie strony.
- Opcjonalnie omija sprawdzanie ograniczeń podczas obciążenia.
- Wykorzystuje zoptymalizowany protokół TDS bulk insert.
- Osiąga przepustowość porównywalną z
bcp.exeiSqlBulkCopy.
Natywne rozszerzenie oparte na języku Rust mssql_py_core obsługuje funkcję zbiorczego kopiowania. Działa poza normalnym potokiem kursora execute().
Podstawowy sposób użycia
Wywołaj bulkcopy() na kursorze, przekazując nazwę tabeli docelowej oraz obiekt iterowalny krotek wierszy lub obiektów Row:
Important
Jeśli tworzysz lub zmieniasz docelową tabelę w tej samej sesji, wywołaj conn.commit() przed bulkcopy(). Protokół kopiowania masowego korzysta z osobnego wewnętrznego kanału do odczytu metadanych tabeli, więc niezatwierdzona zmiana DDL może spowodować impas lub timeout.
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
# Create a temp table for the demo
cursor.execute("""
CREATE TABLE ##BulkDemo (
ID INT,
Name NVARCHAR(50),
Amount MONEY
)
""")
conn.commit()
data = [
(1, "Alice", 50000.00),
(2, "Bob", 60000.00),
(3, "Carol", 55000.00),
]
result = cursor.bulkcopy("##BulkDemo", data)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batch(es)")
print(f"Elapsed: {result['elapsed_time']}")
Wartość zwracana
bulkcopy() Zwraca słownik:
| Klucz | Typ | Opis |
|---|---|---|
rows_copied |
int | Liczba pomyślnie skopiowanych wierszy. |
batch_count |
int | Liczba przetworzonych partii. |
elapsed_time |
float | Czas operacji trwał kilka sekund. |
Sygnatura metody
cursor.bulkcopy(
table_name, # str - target table (can include schema, e.g. "dbo.MyTable")
data, # Iterable[Tuple | Row] - rows to insert
batch_size=0, # int - rows per batch; 0 = server optimal
timeout=30, # int - operation timeout in seconds
column_mappings=None, # List[str] | List[Tuple[int,str]] | None
keep_identity=False, # bool - preserve identity values from source
check_constraints=False, # bool - check constraints during load
table_lock=False, # bool - use table-level lock
keep_nulls=False, # bool - preserve NULLs instead of defaults
fire_triggers=False, # bool - fire INSERT triggers on target
use_internal_transaction=False, # bool - use internal transaction per batch
)
Mapowania kolumn
Domyślnie kolumny bulkcopy() mapowane są według pozycji porządkowej. Każda kolumna danych odpowiada kolumnie tabeli o tym samym indeksie. Użyj parametru column_mappings , aby nadpisać to zachowanie.
Lista nazw kolumn
Każda pozycja na liście odpowiada indeksowi danych źródłowych:
result = cursor.bulkcopy(
"##BulkDemo",
data,
column_mappings=["ID", "Name", "Amount"],
)
Format zaawansowany: jawne mapowanie indeksów
Każda krotka przyjmuje postać (source_index, target_column_name). Użyj tego formatu, aby pominąć lub zmienić kolejność kolumn:
result = cursor.bulkcopy(
"##BulkDemo",
data,
column_mappings=[(0, "ID"), (1, "Name"), (2, "Amount")],
)
Ładuj z plików
Możesz załadować dane z plików CSV i innych formatów plików, przekazując generator do bulkcopy().
Plik CSV
import csv
import io
import mssql_python
# In production, replace io.StringIO with open("data.csv", "r", ...)
csv_data = """ID,Name,Value
1,Widget,9.99
2,Gadget,24.50
3,Gizmo,4.75
"""
def csv_row_generator(file_obj):
"""Generator that yields tuples from a CSV file object."""
reader = csv.reader(file_obj)
next(reader) # Skip header
for row in reader:
if row: # skip blank lines
yield (
int(row[0]), # ID
row[1], # Name
float(row[2]), # Value
)
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##CSVImport (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##CSVImport", csv_row_generator(io.StringIO(csv_data)))
print(f"Imported {result['rows_copied']} rows from CSV")
Duże pliki z grupowaniem
Ustaw batch_size parametr kontrolujący liczbę wierszy, które sterownik wysyła na partię. To podejście sprawdza się dobrze dla dużych plików:
import csv
import io
import mssql_python
# In production, replace io.StringIO with open("large_file.csv", "r", ...)
csv_data = "\n".join(
["ID,Name,Value"] + [f"{i},Item {i},{i * 1.5}" for i in range(1, 201)]
)
def csv_rows(file_obj):
reader = csv.reader(file_obj)
next(reader) # Skip header
for row in reader:
if row:
yield (int(row[0]), row[1], float(row[2]))
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##LargeCSV (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy(
"##LargeCSV",
csv_rows(io.StringIO(csv_data)),
batch_size=50,
)
print(f"Imported {result['rows_copied']} rows in {result['batch_count']} batches")
Load pandas DataFrames
DataFrame ma układ kolumnowy, więc najszybszą ścieżką jest bulkcopy_arrow(), które wykorzystuje tabelę Arrow, którą pandas już potrafi wygenerować.
bulkcopy() przyjmuje krotki reprezentujące wiersze, więc najpierw musisz przekonwertować kolumny na obiekty Pythona.
Przekonwertuj tabelę Arrow do typów kolumn docelowych przed jej załadowaniem.
pyarrow wywnioskowuje float64 dla kolumny numerycznej, której sterownik nie może zmapować na money, decimal ani numeric:
import pandas as pd
import pyarrow as pa
import mssql_python
df = pd.DataFrame({
'ID': [1, 2, 3],
'Name': ['Alice', 'Bob', 'Carol'],
'Amount': [50000.0, 60000.0, 55000.0],
})
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##PandasDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()
target = pa.schema([
pa.field('ID', pa.int32()),
pa.field('Name', pa.string()),
pa.field('Amount', pa.decimal128(19, 4)), # MONEY
])
table = pa.Table.from_pandas(df, preserve_index=False).cast(target)
result = cursor.bulkcopy_arrow("##PandasDemo", table)
Bez rzutowania ładowanie kończy się niepowodzeniem z komunikatem ValueError: Cannot map Arrow column 'Amount' (Float64) to SQL column 'Amount' (Money). Zbuduj rzutowanie za pomocą Table.cast(), zamiast przekazywać schemat do Table.from_pandas(), które nie może bezpośrednio przekonwertować kolumny float na decimal128. W tym przebiegu wartości NaN przyjmują wartość SQL NULL, więc nie musisz ich najpierw zastępować.
Jeśli potrzebujesz ścieżki wiersz-krotki, itertuples() już daje krotki po przejściu name=None:
data = list(df.itertuples(index=False, name=None))
result = cursor.bulkcopy("##PandasDemo", data)
Załaduj dane Apache Arrow
Używa się cursor.bulkcopy_arrow() do ładowania danych Apache Arrow. Ta metoda odczytuje dane bezpośrednio z pamięci Arrow, więc przed jej wywołaniem nie trzeba tworzyć krotek wierszy w Pythonie.
Argument source przyjmuje pyarrow.Table, pyarrow.RecordBatch, pyarrow.RecordBatchReader lub dowolny obiekt, który udostępnia interfejs danych Arrow C. Pozostałe argumenty są takie same jak bulkcopy().
import mssql_python
import pyarrow as pa
conn = mssql_python.connect(connection_string)
# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##ArrowDemo (ID INT, Name NVARCHAR(50), Amount FLOAT)
""")
table = pa.table({
"ID": pa.array([1, 2, 3], type=pa.int32()),
"Name": pa.array(["Alice", "Bob", "Carol"], type=pa.string()),
"Amount": pa.array([50000.0, 60000.0, 55000.0], type=pa.float64()),
})
result = cursor.bulkcopy_arrow("##ArrowDemo", table)
print(f"Copied {result['rows_copied']} rows")
Każdy typ kolumny Arrow musi być zgodny z docelowym typem kolumny SQL. Writer nie konwertuje między rodzinami typów, więc przekazanie kolumny float64 do kolumny pieniężnej podnosi ValueError się zanim zostaną zapisane jakiekolwiek wiersze. Używaj decimal128 dla pieniędzy, kolumn dziesiętnych i liczbowych .
Przekazanie źródła strzałki do bulkcopy() podnosi TypeError i kieruje cię do bulkcopy_arrow().
Więcej informacji o wsparciu dla Arrow, w tym o sposobie przesyłania wyników z jednej tabeli do innej, można znaleźć w artykule Integracja z Apache Arrow.
Obsługa wartości NULL
Przekaż None w dowolnej pozycji kolumny, aby wstawić wartość SQL NULL:
cursor.execute("""
CREATE TABLE ##NullDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()
data = [
(1, "Alice", 50000.00),
(2, "Bob", None), # NULL Amount
(3, None, 55000.00), # NULL Name
]
cursor.bulkcopy("##NullDemo", data)
Kolumny identyfikacyjne
Aby wstawić jawne wartości tożsamościowe, ustaw keep_identity=True:
cursor.execute("""
CREATE TABLE ##IdentDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()
data = [
(100, "Alice", 50000.00),
(200, "Bob", 60000.00),
]
cursor.bulkcopy("##IdentDemo", data, keep_identity=True)
Gdy keep_identity=False (ustawienie domyślne), pomiń kolumnę identity z danych i użyj column_mappings, aby wskazać kolumny inne niż identity.
Opcje kopiowania masowego
| Parameter | Domyślnie | Opis |
|---|---|---|
batch_size |
0 |
Liczba wierszy na wsad.
0 pozwala serwerowi wybrać optymalny rozmiar. |
timeout |
30 |
Limit czasu operacji w sekundach. Dotyczy samej operacji kopiowania masowego, a nie połączenia wewnętrznego. |
keep_identity |
False |
Zachowaj wartości identyfikatorów ze źródłowych danych. |
check_constraints |
False |
Sprawdź ograniczenia tabeli podczas ładowania. |
table_lock |
False |
Załóż blokadę na poziomie tabeli zamiast blokad na poziomie wierszy. |
keep_nulls |
False |
Zachowaj wartości NULL zamiast wstawiać domyślne kolumny. |
fire_triggers |
False |
Spusty ognia INSERT na stole celu. |
use_internal_transaction |
False |
Opakuj każdą partię w transakcji wewnętrznej. |
Note
bulkcopy() otwiera osobne wewnętrzne połączenie z serwerem. To połączenie wewnętrzne dziedziczy limit czasu zapytania kursora: ustaw wartość Connection.timeout na dodatnią przed utworzeniem kursora, a ta sama wartość będzie ograniczać czas próby nawiązania połączenia dla kopiowania zbiorczego. Jeśli limit czasu zapytania kursora wynosi 0, wewnętrzne połączenie korzysta z domyślnego 15-sekundowego limitu czasu nawiązania połączenia. Kursor przyjmuje wartość w momencie utworzenia, więc późniejsza zmiana Connection.timeout nie wpływa na istniejący kursor ani kopię hurtową w trakcie lotu. Zwiększ limit czasu zapytania przed utworzeniem kursora dla wolnych punktów końcowych, punktów końcowych o ograniczonej przepustowości lub o dużych opóźnieniach (na przykład przez VPN lub między regionami).
Zarządzanie błędami
bulkcopy() zgłasza wyjątek, jeśli ładowanie się nie powiedzie, więc opakuj wywołanie w blok try/except, aby przechwycić błędy. Pamiętaj, że bulkcopy() działa na własnym połączeniu wewnętrznym i niezależnie zatwierdza skopiowane wiersze, więc conn.rollback() w głównym połączeniu nie może ich cofnąć. Aby uczynić partię atomową, ustaw use_internal_transaction=True, która owija każdą partię osobną transakcją, która automatycznie cofa się, jeśli partia się nie powiedzie:
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##ImportDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()
data = [
(1, "Alice", 50000.00),
(2, "Bob", 60000.00),
(3, "Carol", 55000.00),
]
try:
result = cursor.bulkcopy("##ImportDemo", data, use_internal_transaction=True)
print(f"Successfully copied {result['rows_copied']} rows")
except (mssql_python.DatabaseError, ValueError) as e:
# bulkcopy() commits on its own connection, so there's nothing to roll back
# here. With use_internal_transaction=True, a failed batch is already rolled
# back on the bulk copy connection.
print(f"Bulk copy failed: {e}")
Aby objąć operację ładowania własną logiką walidacji, skopiuj zbiorczo dane do tabeli przejściowej, a następnie przenieś wiersze do tabeli docelowej za pomocą INSERT ... SELECT w ramach transakcji na głównym połączeniu. To INSERT działa na twoim połączeniu, więc conn.rollback() cofa to, jeśli weryfikacja się nie powiedzie.
Authentication
Kopia masowa korzysta z osobnego wewnętrznego kanału, który wymaga własnego tokena. Sterownik automatycznie obsługuje pozyskiwanie tokenów dla obsługiwanych metod uwierzytelniania.
Zarządzana tożsamość (ActiveDirectoryMSI)
Użyj Authentication=ActiveDirectoryMSI dla tożsamości zarządzanej przypisanej przez system lub przez użytkownika. Ta metoda uwierzytelniania jest zalecana dla usług hostowanych w Azure, takich jak maszyny wirtualne Azure, App Service, Functions i AKS.
import mssql_python
# System-assigned managed identity
conn = mssql_python.connect(
"Server=<server>.database.windows.net;"
"Database=<database>;"
"Authentication=ActiveDirectoryMSI;"
"Encrypt=yes"
)
cursor = conn.cursor()
cursor.execute("CREATE TABLE ##MsiDemo (ID INT, Name NVARCHAR(50))")
conn.commit()
result = cursor.bulkcopy("##MsiDemo", [(1, "Alice"), (2, "Bob")])
print(f"Copied {result['rows_copied']} rows")
W przypadku tożsamości zarządzanej przypisanej przez użytkownika przekaż identyfikator klienta w parametrach połączenia:
conn = mssql_python.connect(
"Server=<server>.database.windows.net;"
"Database=<database>;"
"Authentication=ActiveDirectoryMSI;"
"UID=<client-id>;"
"Encrypt=yes"
)
Nazwa główna usługi (ActiveDirectoryServicePrincipal)
Zastosowanie Authentication=ActiveDirectoryServicePrincipal do uwierzytelniania przez podmiot usługi (dane klienta).
conn = mssql_python.connect(
"Server=<server>.database.windows.net;"
"Database=<database>;"
"Authentication=ActiveDirectoryServicePrincipal;"
"UID=<application-client-id>;"
"PWD=<client-secret>;"
"Encrypt=yes"
)
cursor = conn.cursor()
cursor.execute("CREATE TABLE ##SpDemo (ID INT, Value FLOAT)")
conn.commit()
result = cursor.bulkcopy("##SpDemo", [(1, 1.5), (2, 2.5)])
print(f"Copied {result['rows_copied']} rows")
Domyślny łańcuch poświadczeń (ActiveDirectoryDefault)
ActiveDirectoryDefault Testuje kolejno wielu dostawców poświadczeń, takich jak zmienne środowiskowe, tożsamość obciążenia, tożsamość zarządzana i inne. Działa zarówno dla lokalnego rozwoju, jak i usług hostowanych w Azure bez zmian w kodzie.
Więcej informacji o uwierzytelnianiu można znaleźć w artykule Microsoft Entra authentication.
Porady dotyczące wydajności
Poniższe techniki pomagają zmaksymalizować przepustowość kopiowania masowego.
Zacznij od źródła kolumnowego
bulkcopy()Wymaga iterowalnej liczby krotk wierszowych, więc każda wartość musi istnieć jako obiekt Python zanim kopiowanie się rozpocznie. Gdy dane są już w formacie kolumnowym, bulkcopy_arrow() odczytuje bezpośrednio bufory Arrow i pomija ten krok. Pandas lub Polars DataFrame, plik Parquet oraz wynik cursor.arrow() to wszystkie źródła Arrow. Więcej informacji można znaleźć w artykule Load Apache Arrow data.
Używaj generatorów dla dużych zbiorów danych
Generatory minimalizują zużycie pamięci, ponieważ bulkcopy() akceptują dowolne iterowalne:
def data_generator(count):
"""Generate rows without loading all into memory."""
for i in range(count):
yield (i, f"Item {i}", i * 1.5)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##LargeDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##LargeDemo", data_generator(1000))
Używaj zamków stołowych do szybszych załadunków
Gdy nie ma współbieżnych odczytów, ustaw table_lock=True, aby zmniejszyć narzut związany z blokowaniem podczas dużego początkowego ładowania.
result = cursor.bulkcopy(
"##LargeDemo",
data,
table_lock=True,
batch_size=100000,
)
Wyłącz indeksy podczas ładowania
Tymczasowo wyłącz indeksy nieklastrowane przed ładowaniem zbiorczym i odbuduj je później, aby poprawić wydajność:
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE ##IndexDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
cursor.execute("CREATE NONCLUSTERED INDEX IX_Name ON ##IndexDemo(Name)")
conn.commit()
cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo DISABLE")
conn.commit()
result = cursor.bulkcopy("##IndexDemo", data)
conn.commit()
cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo REBUILD")
conn.commit()
Ładuj tabele równolegle
Otwórz osobne połączenie dla każdej tabeli i uruchamiaj obciążenia równocześnie.
import concurrent.futures
def load_table(table_name, rows):
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute(f"CREATE TABLE {table_name} (ID INT, Name NVARCHAR(50), Value FLOAT)")
conn.commit()
result = cursor.bulkcopy(table_name, rows)
conn.commit()
conn.close()
return result["rows_copied"]
data = [(i, f"Item {i}", i * 1.5) for i in range(100)]
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
futures = [
executor.submit(load_table, "##Load1", data),
executor.submit(load_table, "##Load2", data),
executor.submit(load_table, "##Load3", data),
]
for future in concurrent.futures.as_completed(futures):
print(f"Loaded {future.result()} rows")
Porównanie z alternatywami
Poniższa tabela porównuje kopię masową z innymi metodami wstawiania danych.
| Metoda | Przypadek użycia | Wydajność |
|---|---|---|
cursor.bulkcopy_arrow() |
Duże zbiory danych, które już są kolumnowe. | Najszybszy |
cursor.bulkcopy() |
Duże zbiory danych (ponad 1000 wierszy) pochodzące ze źródeł zorientowanych na wiersze. | Szybko |
cursor.executemany() |
Średnie zbiory danych z parametrami. | Umiarkowane |
cursor.execute() w pętli |
Małe zbiory danych z prostą logiką. | Najwolniejszy |