Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
I denna snabbstart använder du mssql-python-drivrutinen för att kopiera data i bulk mellan databaser. Programmet laddar ned tabeller från ett källdatabasschema till lokala Parquet-filer med Apache Arrow och laddar sedan upp dem till en måldatabas med hjälp av metoden med höga prestanda bulkcopy . Du kan använda det här mönstret för att migrera, replikera eller transformera data mellan SQL Server, Azure SQL Database och SQL Database i Fabric.
Drivrutinen mssql-python kräver inga externa beroenden på Windows-datorer. Drivrutinen installerar allt som behövs med en enda pip installation, så att du kan använda den senaste versionen av drivrutinen för nya skript utan att bryta andra skript som du inte har tid att uppgradera och testa.
mssql-python-dokumentation | mssql-python-källkod | Paket (PyPI) | Uv
Förutsättningar
Python 3.10 eller senare
Om du inte redan har Python installerar du Python-körnings - och pip-pakethanteraren från python.org.
Vill du inte använda din egen miljö? Använd Container- och lokal utveckling för att skapa en reproducerbar devcontainer eller en GitHub Codespaces-miljö.
Visual Studio Code med följande tillägg:
Azure Command-Line Interface (CLI) för lösenordslös autentisering i macOS och Linux.
Om du inte redan har
uvföljer du installationsanvisningarna.En källdatabas med
AdventureWorksLTexempelschemat och en giltig reťazec pripojenia.En destinationsdatabas med en giltig reťazec pripojenia. Användaren måste ha behörighet att skapa och skriva till tabeller. Om du inte har en andra databas kan du använda samma databas och ett annat schema för destinationstabellerna.
Installera operativsystemsspecifika engångsförutsättningar. Windows-användare kan hoppa över detta steg. För fullständiga plattformsdetaljer, se Installera mssql-python.
Skapa en SQL-databas
Skapa eller koppla till en SQL-databas på en av följande plattformar:
Skapa projektet och kör koden
- Skapa ett nytt projekt
- Lägga till beroenden
- Starta Visual Studio Code
- Uppdatera pyproject.toml
- Uppdatera main.py
- Spara anslutningssträngarna
- Använd uv run för att köra skriptet
Skapa ett nytt projekt
Öppna en kommandotolk i utvecklingskatalogen. Om du inte har någon, skapa en ny katalog, till exempel
pythonellerscripts. Undvik mappar på din OneDrive, eftersom synkronisering kan störa hanteringen av din virtuella miljö.Skapa ett nytt projekt med
uv.uv init mssql-python-bcp-qs cd mssql-python-bcp-qs
Lägga till beroenden
Installera paketen mssql-python, python-dotenvoch pyarrow i samma katalog.
uv add mssql-python python-dotenv pyarrow
Öppna Visual Studio Code
Kör följande kommando i samma katalog.
code .
Uppdatera pyproject.toml
pyproject.toml innehåller metadata för projektet. Öppna filen i din favoritredigerare.
Granska innehållet i filen. Det bör likna det här exemplet. Observera Python-versionen och beroendet där
mssql-pythonanvänder>=för att definiera en lägsta version. Om du föredrar en exakt version, ändra>=före versionsnumret till==. De lösta versionerna av varje paket lagras sedan i uv.lock. Låsfilen säkerställer att utvecklare som arbetar med projektet använder konsekventa paketversioner. Det säkerställer också att exakt samma uppsättning paketversioner används när du distribuerar paketet till slutanvändare. Du bör inte redigerauv.lockfilen.[project] name = "mssql-python-bcp-qs" version = "0.1.0" description = "Add your description here" readme = "README.md" requires-python = ">=3.11" dependencies = [ "mssql-python>=1.5.0", "python-dotenv>=1.1.1", "pyarrow>=19.0.0", ]Uppdatera beskrivningen så att den blir mer beskrivande.
description = "Bulk copies data between SQL databases using mssql-python and Apache Arrow"Spara och stäng filen.
Uppdatera main.py
Öppna filen med namnet
main.py. Det bör likna det här exemplet.def main(): print("Hello from mssql-python-bcp-qs!") if __name__ == "__main__": main()Ersätt innehållet i
main.pymed följande kodblock. Varje block bygger på det föregående och bör placeras imain.pyordning.Tips/Råd
Om Visual Studio Code har problem med att lösa paket måste du uppdatera tolken så att den använder den virtuella miljön.
Lägg till importer och konstanter överst i
main.py. Skriptet användermssql_pythonför databasanslutning och hämtning med Arrow,pyarrowochpyarrow.parquetför hantering av kolumnbaserade data och Parquet-fil-I/O,python-dotenvför att läsa in anslutningssträngar från en.env-fil samt ett kompilerat regex-mönster som validerar SQL-identifierare för att förhindra SQL-injektion."""Round-trip: download tables from a source DB/schema to parquet, upload to a destination DB/schema.""" import os import re import time import pyarrow as pa import pyarrow.parquet as pq from dotenv import load_dotenv import mssql_python BATCH_SIZE = 64_000 _SAFE_IDENT = re.compile(r"^[A-Za-z0-9_]+$") def _validate_ident(name: str) -> str: if not _SAFE_IDENT.match(name): raise ValueError(f"Unsafe SQL identifier: {name!r}") return nameLägg till SQL-to-Arrow-typmappningen under importerna. Den här ordlistan översätter SQL Server-kolumntyper till deras Apache Arrow-motsvarigheter så att datafidelitet bevaras vid skrivning till Parquet. Hjälpfunktionerna bygger exakta SQL-typsträngar (till exempel
NVARCHAR(100)ellerDECIMAL(18,2)) frånINFORMATION_SCHEMAmetadata och löser matchande Arrow-typ för varje kolumn. Dessa typer lagras som fältmetadata i Parquet-filerna så att destinationstabellen kan återskapas med exakta kolumndefinitioner._SQL_TO_ARROW = { "bit": pa.bool_(), "tinyint": pa.uint8(), "smallint": pa.int16(), "int": pa.int32(), "bigint": pa.int64(), "float": pa.float64(), "real": pa.float32(), "smallmoney": pa.decimal128(10, 4), "money": pa.decimal128(19, 4), "date": pa.date32(), "datetime": pa.timestamp("us"), "datetime2": pa.timestamp("us"), "smalldatetime": pa.timestamp("s"), "uniqueidentifier": pa.string(), "xml": pa.string(), "image": pa.binary(), "binary": pa.binary(), "varbinary": pa.binary(), "timestamp": pa.binary(), } def _sql_type_str(data_type: str, max_length: int, precision: int, scale: int) -> str: """Build the exact SQL type string from INFORMATION_SCHEMA metadata.""" dt = data_type.lower() if dt in ("char", "varchar", "nchar", "nvarchar", "binary", "varbinary"): length = "MAX" if max_length == -1 else str(max_length) return f"{dt.upper()}({length})" if dt in ("decimal", "numeric"): return f"{dt.upper()}({precision},{scale})" return dt.upper() def _arrow_type(sql_type: str, precision: int, scale: int) -> pa.DataType: sql_type = sql_type.lower() if sql_type in _SQL_TO_ARROW: return _SQL_TO_ARROW[sql_type] if sql_type in ("decimal", "numeric"): return pa.decimal128(precision, scale) if sql_type in ("char", "varchar", "nchar", "nvarchar", "text", "ntext", "sysname"): return pa.string() return pa.string()Lägg till funktionerna för schema introspektion och DDL-generering.
_get_arrow_schemaINFORMATION_SCHEMA.COLUMNSfrågor med hjälp av parametriserade frågor, skapar ett Arrow-schema och lagrar den ursprungliga SQL-typen som fältmetadata så att måltabellen kan återskapas med exakta kolumndefinitioner._create_table_ddlläser dessa metadata tillbaka för att genereraDROP/CREATE TABLEDDL. Typentimestamp(rowversion) mappas om tillVARBINARY(8)eftersom den genereras automatiskt och inte kan infogas direkt.def _get_arrow_schema(cursor, schema_name: str, table_name: str) -> pa.Schema: """Build an Arrow schema from INFORMATION_SCHEMA.COLUMNS. Stores the original SQL type as field metadata so the round-trip CREATE TABLE can reproduce exact column definitions. """ cursor.execute( "SELECT COLUMN_NAME, DATA_TYPE, " "COALESCE(CHARACTER_MAXIMUM_LENGTH, 0), " "COALESCE(NUMERIC_PRECISION, 0), " "COALESCE(NUMERIC_SCALE, 0), " "IS_NULLABLE " "FROM INFORMATION_SCHEMA.COLUMNS " "WHERE TABLE_SCHEMA = ? AND TABLE_NAME = ? " "ORDER BY ORDINAL_POSITION", (schema_name, table_name), ) rows = cursor.fetchall() if not rows: raise ValueError(f"No columns found for {schema_name}.{table_name}") fields = [] for col_name, data_type, max_len, precision, scale, nullable in rows: arrow_t = _arrow_type(data_type, precision, scale) sql_t = _sql_type_str(data_type, max_len, precision, scale) fields.append( pa.field( col_name, arrow_t, nullable=(nullable == "YES"), metadata={"sql_type": sql_t}, ) ) return pa.schema(fields) def _create_table_ddl(target: str, schema: pa.Schema) -> str: """Build DROP/CREATE TABLE DDL from Arrow schema with SQL type metadata.""" col_defs = [] for f in schema: sql_t = f.metadata[b"sql_type"].decode() # timestamp/rowversion is auto-generated and not insertable if sql_t == "TIMESTAMP": sql_t = "VARBINARY(8)" null = "" if f.nullable else " NOT NULL" col_defs.append(f"[{f.name}] {sql_t}{null}") col_defs_str = ",\n ".join(col_defs) return ( f"IF OBJECT_ID('{target}', 'U') IS NOT NULL DROP TABLE {target};\n" f"CREATE TABLE {target} (\n {col_defs_str}\n);" )Lägg till nedladdningsfunktionen.
download_tableanvändercursor.arrow_batch()för att hämta data direkt som Arrow-recordbatchar i drivrutinens C++-lager, vilket gör att man undviker att skapa mellanliggande Python-objekt. Varje batch kastas till det metadata-berikade schemat från_get_arrow_schemaså att de ursprungliga SQL-typerna (till exempelNVARCHAR(100)) bevaras i Parquet-filen. Funktionen använder två separata markörer: en för att läsa kolumnmetadata och en annan för att strömma data.def download_table(conn, schema_name: str, table_name: str, parquet_file: str) -> int: """Download a SQL table to a parquet file. Returns row count (0 if empty).""" _validate_ident(schema_name) _validate_ident(table_name) source = f"{schema_name}.[{table_name}]" with conn.cursor() as cursor: schema = _get_arrow_schema(cursor, schema_name, table_name) row_count = 0 t0 = time.perf_counter() with conn.cursor() as cursor: cursor.execute(f"SELECT * FROM {source}") writer = None try: while True: batch = cursor.arrow_batch(BATCH_SIZE) if batch.num_rows == 0: break # Cast to the schema to preserve SQL type metadata in Parquet arrays = [ batch.column(i).cast(schema.field(i).type) for i in range(batch.num_columns) ] batch = pa.record_batch(arrays, schema=schema) if writer is None: writer = pq.ParquetWriter(parquet_file, schema) writer.write_batch(batch) row_count += batch.num_rows finally: if writer is not None: writer.close() if row_count == 0: return 0 elapsed = time.perf_counter() - t0 rate = f"{int(row_count / elapsed):,} rows/sec" if elapsed > 0 else "n/a" print( f"{schema_name}.{table_name} → {parquet_file}: {row_count:,} rows downloaded " f"in {elapsed:.2f}s ({rate})" ) return row_countLägg till enrichment hook.
enrich_parquetär en platshållare där du kan lägga till transformeringar, härledda kolumner eller kopplingar till data innan de laddas upp. I den här snabbstarten är det en no-op som returnerar filsökvägen oförändrad.def enrich_parquet(parquet_file: str) -> str: """Enrich a parquet file before upload. Returns the (possibly new) file path.""" # TODO: add transformations, derived columns, or joins print(f"Enriching {parquet_file} (no-op)") return parquet_fileLägg till uppladdningsfunktionen.
upload_parquetläser Arrow-schemat från Parquet-filen, genererar och körDROP/CREATE TABLEDDL för att förbereda målet, läser sedan filen i batchar och anroparcursor.bulkcopy()för snabb bearbetning vid bulk-infogning. Alternativettable_lock=Trueförbättrar genomströmningen genom att minimera låskonkurrensen. Efter att uppladdningen är klar kör funktionen enSELECT COUNT(*)och ger ett fel om antalet destinationsrader inte matchar antalet uppladdade rader.def upload_parquet(conn, parquet_file: str, target: str) -> int: """Upload a parquet file into a SQL table via BCP. Returns row count.""" # ── Create target table from parquet schema ── pf_schema = pq.read_schema(parquet_file) with conn.cursor() as cursor: cursor.execute(_create_table_ddl(target, pf_schema)) conn.commit() # ── Bulk insert ── uploaded = 0 t0 = time.perf_counter() with pq.ParquetFile(parquet_file) as pf: with conn.cursor() as cursor: for batch in pf.iter_batches(batch_size=BATCH_SIZE): rows = zip(*(col.to_pylist() for col in batch.columns)) cursor.bulkcopy( target, rows, batch_size=BATCH_SIZE, table_lock=True, timeout=3600, ) uploaded += batch.num_rows elapsed = time.perf_counter() - t0 # ── Verify ── with conn.cursor() as cursor: cursor.execute(f"SELECT COUNT(*) FROM {target}") count = cursor.fetchone()[0] if count != uploaded: raise ValueError( f"Row count mismatch for {target}: uploaded {uploaded:,}, destination has {count:,}" ) rate = f"{int(uploaded / elapsed):,} rows/sec" if elapsed > 0 else "n/a" print( f"{parquet_file} → {target}: {uploaded:,} rows uploaded " f"in {elapsed:.2f}s ({rate}) " f"| destination rows: {count:,}" ) return uploadedLägg till orkestreringsfunktionen.
transfer_tablesbinder ihop de tre faserna. Den ansluter till källdatabasen, upptäcker alla bastabeller i det angivna schemat viaINFORMATION_SCHEMA.TABLES, laddar ned var och en till en lokal Parquet-fil, kör enrichment-hooken, ansluter därefter till destinationsdatabasen och laddar upp varje fil.def transfer_tables( source_conn_str: str, dest_conn_str: str, source_schema: str, dest_schema: str, ) -> None: """Download all tables from source DB/schema to parquet, upload to dest DB/schema.""" _validate_ident(source_schema) _validate_ident(dest_schema) parquet_dir = source_schema os.makedirs(parquet_dir, exist_ok=True) # ── Download from source ── with mssql_python.connect(source_conn_str) as src_conn: with src_conn.cursor() as cursor: cursor.execute( "SELECT TABLE_NAME FROM INFORMATION_SCHEMA.TABLES " "WHERE TABLE_SCHEMA = ? AND TABLE_TYPE = 'BASE TABLE' " "ORDER BY TABLE_NAME", (source_schema,), ) tables = [row[0] for row in cursor.fetchall()] print(f"Found {len(tables)} {source_schema} tables: {', '.join(tables)}\n") parquet_files = [] for table_name in tables: parquet_file = os.path.join(parquet_dir, f"{table_name}.parquet") row_count = download_table(src_conn, source_schema, table_name, parquet_file) if row_count == 0: print(f"{source_schema}.{table_name}: empty, skipping") else: parquet_files.append((table_name, parquet_file)) # ── Enrich parquet files ── enriched = [] for table_name, parquet_file in parquet_files: enriched.append((table_name, enrich_parquet(parquet_file))) # ── Upload to destination ── with mssql_python.connect(dest_conn_str) as dest_conn: for table_name, parquet_file in enriched: target = f"{dest_schema}.[{table_name}]" upload_parquet(dest_conn, parquet_file, target)Lägg slutligen till
mainstartpunkten. Den läser in.envfilen, anropartransfer_tablesmed käll- och målanslutningssträngarna och skriver ut den totala förflutna tiden.def main(): load_dotenv() t_start = time.perf_counter() transfer_tables( source_conn_str=os.environ["SOURCE_CONNECTION_STRING"], dest_conn_str=os.environ["DEST_CONNECTION_STRING"], source_schema="SalesLT", dest_schema="dbo", ) print(f"Total: {time.perf_counter() - t_start:.2f}s") if __name__ == "__main__": main()Spara och stäng
main.py.
Spara anslutningssträngarna
.gitignoreÖppna filen och lägg till ett undantag för.envfiler. Filen bör likna det här exemplet. Se till att spara och stänga den när du är klar.# Python-generated files __pycache__/ *.py[oc] build/ dist/ wheels/ *.egg-info # Virtual environments .venv # Connection strings and secrets .envI den aktuella katalogen skapar du en ny fil med namnet
.env..envLägg till poster för käll- och målanslutningssträngarna i filen. Ersätt platshållarvärdena med dina faktiska server- och databasnamn.SOURCE_CONNECTION_STRING="Server=<source_server_name>;Database=<source_database_name>;Encrypt=yes;TrustServerCertificate=no;Authentication=ActiveDirectoryInteractive" DEST_CONNECTION_STRING="Server=<dest_server_name>;Database=<dest_database_name>;Encrypt=yes;TrustServerCertificate=no;Authentication=ActiveDirectoryInteractive"Tips/Råd
Anslutningssträngen som används här beror till stor del på vilken typ av SQL-databas du ansluter till. Om du ansluter till en Azure SQL Database eller en SQL-databas i Fabric använder du ODBC-anslutningssträngen från fliken Anslutningssträngar. Du kan behöva justera autentiseringstypen beroende på ditt scenario. Mer information om anslutningssträngar och deras syntax finns i referens för anslutningssträngssyntax.
Tips/Råd
På macOS fungerar både ActiveDirectoryInteractive och ActiveDirectoryDefault för Microsoft Entra-autentisering.
ActiveDirectoryInteractive uppmanar dig att logga in varje gång du kör skriptet. För att undvika upprepade inloggningspromptar, logga in en gång via Azure CLI genom att köra az login, och använd ActiveDirectoryDefaultsedan , vilket återanvänder den cachade legitimationen.
Använd uv run för att köra skriptet
Kör följande kommando i terminalfönstret från tidigare eller ett nytt terminalfönster som är öppet till samma katalog.
uv run main.pyHär är de förväntade utdata när skriptet är klart.
Found 12 SalesLT tables: Address, Customer, CustomerAddress, ... SalesLT.Address → SalesLT/Address.parquet: 450 rows downloaded in 0.15s (3,000 rows/sec) ... SalesLT/Address.parquet → dbo.[Address]: 450 rows uploaded in 0.10s (4,500 rows/sec) | verified: 450 ... Total: 2.35sAnslut dig till destinationsdatabasen genom att använda MSSQL-tillägget för VS Code och verifiera att tabellerna och datan har skapats korrekt.
Om du vill distribuera skriptet till en annan dator kopierar du alla filer förutom
.venvmappen till den andra datorn. Den virtuella miljön återskapas med den första körningen.
Så här fungerar koden
Programmet utför en fullständig dataöverföring tur och retur i tre faser:
-
Ladda ned: Ansluter till källdatabasen, läser kolumnmetadata från
INFORMATION_SCHEMA.COLUMNS, skapar ett Apache Arrow-schema och laddar sedan ned varje tabell till en lokal Parquet-fil. -
Berika (valfritt): Ger en punkt (
enrich_parquet) där du kan lägga till transformeringar, härledda kolumner eller sammanfogningar innan du laddar upp. -
Ladda upp: Läser varje Parquet-fil i batchar, återskapar tabellen i måldatabasen med DDL som genererats från Arrow-schemametadata och använder
cursor.bulkcopy()sedan för massinfogning med höga prestanda.
Nästa steg
Använd dessa artiklar för att fortsätta bygga:
- Masskopiering för avancerade bulkkopieringsmönster inklusive kolumnmappningar, batchstorlek och felhantering.
- Dataladdnings- och rörelsemönster för strategier för att ladda data från filer, API:er och andra databaser.
- Prestandaoptimering för att optimera genomströmningen för stora dataoperationer.