Quickstart: Bulksgewijs kopiëren met het mssql-python-stuurprogramma voor Python

In deze quickstart gebruikt u het mssql-python stuurprogramma om gegevens bulksgewijs tussen databases te kopiëren. De toepassing downloadt tabellen van een brondatabaseschema naar lokale Parquet-bestanden met behulp van Apache Arrow en uploadt ze vervolgens naar een doeldatabase met behulp van de krachtige bulkcopy methode. U kunt dit patroon gebruiken om gegevens te migreren, repliceren of transformeren tussen SQL Server, Azure SQL Database en SQL Database in Fabric.

Het mssql-python stuurprogramma vereist geen externe afhankelijkheden op Windows-computers. Het stuurprogramma installeert alles wat het nodig heeft met één pip installatie, zodat u de nieuwste versie van het stuurprogramma voor nieuwe scripts kunt gebruiken zonder dat andere scripts die u niet hoeft te upgraden en te testen, worden onderbroken.

documentatie | mssql-python-broncode | Pakket (PyPI) | Uv

Vereiste voorwaarden

  • Python 3.10 of hoger

  • Als u Nog geen Python hebt, installeert u de Python-runtime en pip-pakketbeheer vanuit python.org.

  • Wilt u niet uw eigen omgeving gebruiken? Volg Container- en lokale ontwikkeling om een reproduceerbare devcontainer of GitHub Codespaces-omgeving te creëren.

  • Visual Studio Code met de volgende extensies:

  • Azure Command-Line Interface (CLI) voor verificatie zonder wachtwoord in macOS en Linux.

  • Als u dat nog niet hebt uv, volgt u de installatie-instructies.

  • Een brondatabase met het AdventureWorksLT voorbeeldschema en een geldige verbindingsreeks.

  • Een bestemmingsdatabase met een geldige verbindingsreeks. De gebruiker moet gemachtigd zijn om tabellen te maken en naar tabellen te schrijven. Als je geen tweede database hebt, kun je dezelfde database en een ander schema gebruiken voor de bestemmingstabellen.

Installeer eenmalige vereisten voor het besturingssysteem. Windows-gebruikers kunnen deze stap overslaan. Voor volledige platformdetails, zie Install mssql-python.

apk add libtool krb5-libs krb5-dev

Een SQL-database maken

Maak een SQL-database aan of maak verbinding met een van de volgende platforms:

Het project maken en de code uitvoeren

  1. Een nieuw project maken
  2. Afhankelijkheden toevoegen
  3. Visual Studio Code starten
  4. Pyproject.toml bijwerken
  5. Main.py bijwerken
  6. De verbindingsreeksen opslaan
  7. Uv-uitvoering gebruiken om het script uit te voeren

Een nieuw project maken

  1. Open een opdrachtprompt in uw ontwikkelingsmap. Als je er geen hebt, maak dan een nieuwe map aan, zoals python of scripts. Vermijd mappen op je OneDrive, want synchronisatie kan het beheer van je virtuele omgeving verstoren.

  2. Maak een nieuw project met uv.

    uv init mssql-python-bcp-qs
    cd mssql-python-bcp-qs
    

Afhankelijkheden toevoegen

Installeer in dezelfde map de mssql-python, python-dotenven pyarrow pakketten.

uv add mssql-python python-dotenv pyarrow

Visual Studio Code starten

Voer in dezelfde map de volgende opdracht uit.

code .

Pyproject.toml bijwerken

  1. Het pyproject.toml bevat de metagegevens voor uw project. Open het bestand in uw favoriete editor.

  2. Controleer de inhoud van het bestand. Deze moet vergelijkbaar zijn met dit voorbeeld. Noteer de Python-versie en -afhankelijkheid voor mssql-python gebruik >= om een minimale versie te definiëren. Als u de voorkeur geeft aan een exacte versie, wijzigt u de >= voor het versienummer in ==. De opgeloste versies van elk pakket worden vervolgens opgeslagen in de uv.lock. Het lockfile zorgt ervoor dat ontwikkelaars die aan het project werken, consistente pakketversies gebruiken. Het zorgt er ook voor dat dezelfde set pakketversies wordt gebruikt bij het distribueren van uw pakket naar eindgebruikers. U moet het uv.lock bestand niet bewerken.

    [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",
    ]
    
  3. Werk de beschrijving bij zodat deze meer beschrijvend is.

    description = "Bulk copies data between SQL databases using mssql-python and Apache Arrow"
    
  4. Sla het bestand op en sluit het.

Main.py bijwerken

  1. Open het bestand met de naam main.py. Deze moet vergelijkbaar zijn met dit voorbeeld.

    def main():
        print("Hello from mssql-python-bcp-qs!")
    
    if __name__ == "__main__":
        main()
    
  2. Vervang de inhoud van main.py met de volgende codeblokken. Elk blok bouwt voort op de vorige en moet op main.py volgorde worden geplaatst.

    Aanbeveling

    Als Visual Studio Code problemen ondervindt bij het oplossen van pakketten, moet u de interpreter bijwerken om de virtuele omgeving te gebruiken.

  3. Voeg bovenaan main.pyde import- en constanten toe. Het script gebruikt mssql_python voor databaseconnectiviteit en Arrow-ophaal, pyarrow en pyarrow.parquet voor kolomgegevensverwerking en Parquet-bestand-I/O, python-dotenv voor het laden van verbindingsstrings uit een .env bestand, en een gecompileerd regex-patroon dat SQL-identificaties valideert om injectie te voorkomen.

    """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 name
    
  4. Voeg onder de imports de SQL-naar-Arrow type mapping toe. Deze woordenlijst vertaalt SQL Server-kolomtypen naar hun Apache Arrow-equivalenten, zodat de betrouwbaarheid van gegevens behouden blijft bij het schrijven naar Parquet. De helperfuncties bouwen exacte SQL-typestrings (bijvoorbeeld NVARCHAR(100) of DECIMAL(18,2)) uit INFORMATION_SCHEMA metadata en lossen het overeenkomende Arrow-type voor elke kolom op. Deze typen worden opgeslagen als veldmetadata in de Parquet-bestanden zodat de bestemmingstabel kan worden nagemaakt met exacte kolomdefinities.

    _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()
    
  5. Voeg de functies schema-introspectie en DDL-generatie toe. _get_arrow_schema INFORMATION_SCHEMA.COLUMNS-query’s met behulp van geparameteriseerde query's, bouwt een Arrow schema en slaat het oorspronkelijke SQL-type op als veldmetagegevens, zodat de doeltabel opnieuw kan worden gemaakt met exacte kolomdefinities. _create_table_ddl leest die metagegevens terug om DDL te genereren DROP/CREATE TABLE . Het timestamp-type (rowversion) wordt opnieuw gemapt naar VARBINARY(8) omdat het automatisch wordt gegenereerd en niet kan worden ingevoegd.

    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);"
        )
    
  6. Voeg de downloadfunctie toe. download_tablegebruikt cursor.arrow_batch() om data direct op te halen als Arrow-recordbatches in de C++-laag van de driver, waardoor het creëren van tussenliggende Python-objecten voorkomt. Elke batch wordt vanuit het metadata-verrijkte schema _get_arrow_schema gecast, zodat de originele SQL-types (bijvoorbeeld NVARCHAR(100)) behouden blijven in het Parquet-bestand. De functie maakt gebruik van twee afzonderlijke cursors: een voor het lezen van kolommetagegevens en een andere om de gegevens te streamen.

    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_count
    
  7. Voeg de verrijkingshook toe. enrich_parquet is een tijdelijke aanduiding waar u transformaties, afgeleide kolommen of joins kunt toevoegen aan gegevens voordat deze worden geüpload. In deze snelstartgids is er sprake van een no-op die het bestandspad ongewijzigd teruggeeft.

    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_file
    
  8. Voeg de uploadfunctie toe. upload_parquet leest het Arrow-schema uit het Parquet-bestand, genereert en voert DROP/CREATE TABLE DDL uit om de bestemming voor te bereiden. Vervolgens wordt het bestand in batches gelezen en wordt cursor.bulkcopy() aangeroepen voor een bulkinvoer met hoge prestaties. De table_lock=True optie verbetert de doorvoer door vergrendelingsconflicten te minimaliseren. Nadat de upload is voltooid, voert de functie een SELECT COUNT(*) uit en geeft een foutmelding als het aantal bestemmingsrijen niet overeenkomt met het aantal geüploade rijen.

    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 uploaded
    
  9. Voeg de orchestratiefunctie toe. transfer_tables koppelt de drie fasen aan elkaar. Het maakt verbinding met de brondatabase, ontdekt alle basistabellen in het gegeven schema via INFORMATION_SCHEMA.TABLES, downloadt elke tabel naar een lokaal Parquet-bestand, voert de verrijkingshaak uit, verbindt zich vervolgens met de bestemmingsdatabase en uploadt elk bestand.

    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)
    
  10. Voeg tot slot het main toegangspunt toe. Het laadt het bestand .env, roept transfer_tables aan met de verbindingsreeksen voor de bron en het doel, en drukt de totale verstreken tijd af.

    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()
    
  11. Opslaan en sluiten main.py.

De verbindingsreeksen opslaan

  1. Open het .gitignore bestand en voeg een uitsluiting toe voor .env bestanden. Het bestand moet er ongeveer uitzien als in dit voorbeeld. Zorg ervoor dat u deze opslaat en sluit wanneer u klaar bent.

    # Python-generated files
    __pycache__/
    *.py[oc]
    build/
    dist/
    wheels/
    *.egg-info
    
    # Virtual environments
    .venv
    
    # Connection strings and secrets
    .env
    
  2. Maak in de huidige map een nieuw bestand met de naam .env.

  3. Voeg in het .env bestand vermeldingen toe voor uw bron- en doelverbindingsreeksen. Vervang de tijdelijke aanduidingen door de werkelijke server- en databasenamen.

    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"
    

    Aanbeveling

    De hier gebruikte verbindingsreeks is grotendeels afhankelijk van het type SQL-database waarmee u verbinding maakt. Als u verbinding maakt met een Azure SQL Database of een SQL-database in Fabric, gebruikt u de ODBC-verbindingsreeks op het tabblad Verbindingsreeksen. Mogelijk moet u het verificatietype aanpassen, afhankelijk van uw scenario. Zie de naslaginformatie over de syntaxis van de verbindingsreeks voor meer informatie over verbindingsreeksen en de bijbehorende syntaxis.

Aanbeveling

In macOS werken beide ActiveDirectoryInteractive en ActiveDirectoryDefault voor Microsoft Entra-verificatie. ActiveDirectoryInteractive u wordt gevraagd u aan te melden telkens wanneer u het script uitvoert. Om herhaalde aanmeldingsprompts te voorkomen, meld je je één keer aan via de Azure CLI door az login uit te voeren en gebruik vervolgens ActiveDirectoryDefault, waarmee de in de cache opgeslagen referentiegegevens opnieuw worden gebruikt.

Gebruik uv run om het script uit te voeren

  1. Voer in het terminalvenster van vóór, of een nieuw terminalvenster dat is geopend in dezelfde map, de volgende opdracht uit.

     uv run main.py
    

    Dit is de verwachte uitvoer wanneer het script is voltooid.

    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.35s
    
  2. Maak verbinding met de bestemmingsdatabase door gebruik te maken van de MSSQL-extensie voor VS Code, en verifieer dat de tabellen en gegevens succesvol zijn aangemaakt.

  3. Als u uw script op een andere computer wilt implementeren, kopieert u alle bestanden, met uitzondering van de .venv map naar de andere computer. De virtuele omgeving wordt opnieuw gemaakt met de eerste uitvoering.

Hoe de code werkt

De toepassing voert een volledige retourgegevensoverdracht uit in drie fasen:

  1. Download: Maakt verbinding met de brondatabase, leest kolommetagegevens van INFORMATION_SCHEMA.COLUMNS, bouwt een Apache Arrow-schema en downloadt vervolgens elke tabel in een lokaal Parquet-bestand.
  2. Verrijken (optioneel): Biedt een haak (enrich_parquet) waar u transformaties, afgeleide kolommen of joins kunt toevoegen voordat u uploadt.
  3. Uploaden: leest elk Parquet-bestand in batches, maakt de tabel opnieuw in de doeldatabase met behulp van DDL die is gegenereerd op basis van metagegevens van het pijlschema en gebruikt cursor.bulkcopy() vervolgens voor bulksgewijs invoegen met hoge prestaties.

Volgende stappen 

Gebruik deze artikelen om verder te bouwen:

  • Bulk copy voor geavanceerde bulk copy patronen, waaronder kolommappings, batch-dimensionering en foutafhandeling.
  • Data-laad- en bewegingspatronen voor strategieën om data te laden uit bestanden, API's en andere databases.
  • Prestatie-tuning om de doorvoer te optimaliseren voor grote databewerkingen.