Apache Spark-Connector: SQL Server und Azure SQL

Der Apache Spark-Connector für SQL Server und Azure SQL ist ein leistungsstarker Connector, den Sie verwenden können, um Transaktionsdaten in Big Data Analytics einzuschließen und Ergebnisse für Ad-hoc-Abfragen oder Berichte beizubehalten. Mithilfe des Connectors können Sie jede sql-Datenbank, lokal oder in der Cloud als Eingabedatenquelle oder Ausgabedatensenke für Spark-Aufträge verwenden.

Hinweis

Dieser Stecker wird nicht mehr gepflegt. Das Upstream-Projekt wurde im Februar 2025 archiviert, und dieser Artikel wird nur zu Archivzwecken aufbewahrt. Für neue Arbeiten gegen SQL Server oder Azure SQL verwenden Sie die integrierte JDBC-Datenquelle Apache Spark mit dem Microsoft JDBC Driver für SQL Server.

Diese Bibliothek enthält den Quellcode für den Apache Spark Connector für SQL Server- und Azure SQL-Plattformen.

Apache Spark ist eine vereinheitlichte Engine zur Verarbeitung von umfangreichen Daten.

Importiere den Connector mit der Koordinate, die zu deiner Spark-Version passt:

Verbinder Maven-Koordinate
Spark 2.4.x-kompatibler Anschluss com.microsoft.azure:spark-mssql-connector:1.0.2
Spark 3.0.x-kompatibles Verbindungselement com.microsoft.azure:spark-mssql-connector_2.12:1.1.0
Spark 3.1.x-kompatibler Anschluss com.microsoft.azure:spark-mssql-connector_2.12:1.2.0
Spark 3.3.x-kompatibler Anschluss (Beta) com.microsoft.azure:spark-mssql-connector_2.12:1.3.0-BETA
Spark 3.4.x-kompatibler Anschluss (Beta) Keine Maven-Koordinate. Laden Sie das spark-mssql-connector_2.12-1.4.0-BETA.jar Asset von der GitHub-Veröffentlichungsseite herunter.

Die Spark 3.3.x- und Spark 3.4.x-Connectoren sind Beta-Versionen, und keiner von beiden wurde vor der Archivierung des Projekts allgemein verfügbar. Die Version 1.3.0-BETA ist nicht kompatibel mit dem Microsoft JDBC-Treiber für SQL Server 7.0.1, Spark 2.4 oder Spark 3.0. Die Upstream-Release Notes von 1.3.0 empfehlen, die Version 1.3.0zu importieren, aber diese Version wurde nie veröffentlicht. Verwenden Sie stattdessen 1.3.0-BETA.

Sie können den Connector auch aus der Quelle erstellen oder den JAR aus dem Abschnitt "Release" in GitHub herunterladen. Quellcode und Versionsgeschichte finden Sie im archivierten SQL Spark Connector GitHub-Repository.

Unterstützte Funktionen

  • Unterstützung für alle Spark-Bindungen (Scala, Python, R)
  • Unterstützung der Standardauthentifizierung und der Active Directory (AD)-Schlüsseltabelle
  • Geänderte Unterstützung für dataframe-Schreibvorgänge
  • Unterstützung von Schreibvorgängen in SQL Server-Einzelinstanzen und Datenpools in Big Data-Clustern von SQL Server
  • Zuverlässige Connector-Unterstützung für den SQL Server Single Instance
Komponente Unterstützte Versionen
Apache Spark 2.4.x, 3.0.x, 3.1.x, 3.3.x (Beta), 3.4.x (Beta)
Scala 2.11, 2.12
Microsoft JDBC-Treiber für SQL Server 8,4
Microsoft SQL Server SQL Server 2008 oder höher
Azure SQL-Datenbanken Unterstützt

Unterstützte Optionen

Der Apache Spark Connector für SQL Server und Azure SQL unterstützt die im SQL DataSource JDBC-Artikel definierten Optionen.

Darüber hinaus unterstützt der Connector die folgenden Optionen:

Option Standard BESCHREIBUNG
reliabilityLevel BEST_EFFORT BEST_EFFORT oder NO_DUPLICATES. NO_DUPLICATES implementiert ein zuverlässiges Einfügen in Ausführungs-Neustartszenarien
dataPoolDataSource none none impliziert, dass der Wert nicht festgelegt ist und der Connector in eine einzelne SQL Server-Instanz schreiben soll. Legen Sie diesen Wert auf den Namen der Datenquelle fest, um eine Datenpooltabelle in Big Data-Cluster zu schreiben.
isolationLevel READ_COMMITTED Angeben der Isolationsstufe
tableLock false Implementiert ein Einfügen mit TABLOCK Option zur Verbesserung der Schreibgeschwindigkeit.
schemaCheckEnabled true Deaktiviert die strenge Datenrahmen- und SQL-Tabellenschemaüberprüfung, wenn diese auf "false" festgelegt ist.

Legen Sie andere Massenkopieroptionen für den dataframe fest. Der Konnektor übergibt diese Optionen an bulkcopy-APIs beim Schreiben.

Leistungsvergleich

Apache Spark Connector für SQL Server und Azure SQL ist bis zu 15x schneller als generischer JDBC-Connector für die Datenübertragung in SQL Server. Die Leistungsmerkmale variieren je nach Typ, Datenvolumen, verwendeten Optionen und können Abweichungen zwischen den einzelnen Ausführungen darstellen. Die folgenden Leistungsergebnisse sind die Zeitaufwand zum Überschreiben einer SQL-Tabelle mit 143,9M Zeilen in einem Spark dataframe. Der Spark dataframe wird erstellt, indem store_sales die HDFS-Tabelle gelesen wird, die mit Spark TPCDS Benchmark generiert wird. Die Zeit zum Lesen von store_sales in dataframe wird ausgeschlossen. Die Ergebnisse werden über drei Ausführungen gemittelt.

Verbindertyp Optionen BESCHREIBUNG Zeit zum Schreiben
JDBCConnector Standard Generischer JDBC-Connector mit Standardoptionen 1.385 Sekunden
sql-spark-connector BEST_EFFORT Optimale Leistung sql-spark-connector mit Standardoptionen 580 Sekunden
sql-spark-connector NO_DUPLICATES Zuverlässig sql-spark-connector 709 Sekunden
sql-spark-connector BEST_EFFORT + tabLock=true Bestmögliche sql-spark-connector-Leistung mit aktivierter Tabellensperre 72 Sekunden
sql-spark-connector NO_DUPLICATES + tabLock=true Zuverlässige sql-spark-connector-Leistung mit aktivierter Tabellensperre 198 Sekunden

Konfiguration

  • Spark-Konfiguration: num_executors = 20, executor_memory = '1664 MB', executor_cores = 2
  • Data Gen config: scale_factor=50, partitioned_tables=true
  • Datendatei store_sales mit der Anzahl der Zeilen 143.997.590

Umwelt

  • SQL Server Big Data Cluster CU5
  • master + 6 Knoten
  • Jeder Knoten ein Gen-5-Server mit 512 GB RAM, 4 TB NVM pro Knoten und 10-GBit/s NIC

Häufig auftretende Probleme

java.lang.NoClassDefFoundError: com/microsoft/aad/adal4j/AuthenticationException

Dieser Fehler tritt auf, wenn Sie eine ältere Version des mssql Treibers in Ihrer Hadoop-Umgebung verwenden. Der Connector enthält jetzt diesen Treiber. Wenn Sie zuvor den Azure SQL Connector verwendet und Treiber manuell auf Ihrem Cluster für die Kompatibilität der Microsoft Entra-Authentifizierung installiert haben, entfernen Sie diese Treiber.

So beheben Sie den Fehler:

  1. Wenn Sie eine generische Hadoop-Umgebung verwenden, überprüfen und entfernen Sie den mssql JAR mit dem folgenden Befehl: rm $HADOOP_HOME/share/hadoop/yarn/lib/mssql-jdbc-6.2.1.jre7.jar Wenn Sie Databricks verwenden, fügen Sie ein globales Oder Cluster-Init-Skript hinzu, um alte Versionen des mssql Treibers aus dem /databricks/jars Ordner zu entfernen, oder fügen Sie diese Zeile zu einem vorhandenen Skript hinzu: rm /databricks/jars/*mssql*

  2. Fügen Sie die adal4j und mssql Pakete hinzu. Beispielsweise können Sie Maven verwenden, aber jede Möglichkeit sollte funktionieren.

    Vorsicht

    Installieren Sie den SQL Spark Connector nicht auf diese Weise.

  3. Fügen Sie der Verbindungskonfiguration die Treiberklasse hinzu. Beispiel:

    connectionProperties = {
      `Driver`: `com.microsoft.sqlserver.jdbc.SQLServerDriver`
    }`
    

Weitere Informationen finden Sie in der Lösung für https://github.com/microsoft/sql-spark-connector/issues/26.

Get started

Der Apache Spark Connector für SQL Server und Azure SQL basiert auf der Spark DataSourceV1-API und der SQL Server-Massen-API. Es verwendet die gleiche Schnittstelle wie der integrierte JDBC Spark-SQL-Konnektor. Mithilfe dieser Integration können Sie den Connector problemlos integrieren und Ihre vorhandenen Spark-Aufträge migrieren, indem Sie den Formatparameter mit com.microsoft.sqlserver.jdbc.sparkaktualisieren.

Um den Connector in Ihre Projekte einzuschließen, laden Sie dieses Repository herunter, und erstellen Sie den JAR mit SBT.

Schreiben in eine neue SQL-Tabelle

Vorsicht

Der overwrite Modus legt die Tabelle zuerst ab, wenn sie bereits in der Datenbank vorhanden ist. Verwenden Sie diese Option mit Bedacht, um unerwartete Datenverluste zu vermeiden.

Wenn Sie den Modus overwrite ohne die Option truncate beim Erneuten Erstellen der Tabelle verwenden, entfernt der Vorgang Indizes. Außerdem ändert sich eine Columnstore-Tabelle in eine Heap-Tabelle. Um vorhandene Indizes beizubehalten, legen Sie die truncate Option auf true. Beispiel: .option("truncate","true").

server_name = "jdbc:sqlserver://{SERVER_ADDR}"
database_name = "database_name"
url = server_name + ";" + "databaseName=" + database_name + ";"

table_name = "table_name"
username = "username"
password = "password123!#" # Please specify password here

try:
  df.write \
    .format("com.microsoft.sqlserver.jdbc.spark") \
    .mode("overwrite") \
    .option("url", url) \
    .option("dbtable", table_name) \
    .option("user", username) \
    .option("password", password) \
    .save()
except ValueError as error :
    print("Connector write failed", error)

Anfügen an eine SQL-Tabelle

try:
  df.write \
    .format("com.microsoft.sqlserver.jdbc.spark") \
    .mode("append") \
    .option("url", url) \
    .option("dbtable", table_name) \
    .option("user", username) \
    .option("password", password) \
    .save()
except ValueError as error :
    print("Connector write failed", error)

Angeben der Isolationsstufe

Dieser Connector verwendet standardmäßig die READ_COMMITTED Isolationsstufe, wenn er Daten als Massendaten in die Datenbank einfügt. Verwenden Sie die mssqlIsolationLevel Option, um die Isolationsebene außer Kraft zu setzen:

    .option("mssqlIsolationLevel", "READ_UNCOMMITTED") \

Aus SQL-Tabelle lesen

jdbcDF = spark.read \
        .format("com.microsoft.sqlserver.jdbc.spark") \
        .option("url", url) \
        .option("dbtable", table_name) \
        .option("user", username) \
        .option("password", password).load()

Microsoft Entra-Authentifizierung

Python-Beispiel mit Service-Principal

context = adal.AuthenticationContext(authority)
token = context.acquire_token_with_client_credentials(resource_app_id_url, service_principal_id, service_principal_secret)
access_token = token["accessToken"]

jdbc_db = spark.read \
        .format("com.microsoft.sqlserver.jdbc.spark") \
        .option("url", url) \
        .option("dbtable", table_name) \
        .option("accessToken", access_token) \
        .option("encrypt", "true") \
        .option("hostNameInCertificate", "*.database.windows.net") \
        .load()

Python-Beispiel mit Active Directory-Kennwort

jdbc_df = spark.read \
        .format("com.microsoft.sqlserver.jdbc.spark") \
        .option("url", url) \
        .option("dbtable", table_name) \
        .option("authentication", "ActiveDirectoryPassword") \
        .option("user", user_name) \
        .option("password", password) \
        .option("encrypt", "true") \
        .option("hostNameInCertificate", "*.database.windows.net") \
        .load()

Um sich mithilfe von Active Directory zu authentifizieren, installieren Sie die erforderliche Abhängigkeit.

Bei Verwendung von ActiveDirectoryPassword sollte der user-Wert im UPN-Format vorliegen, z. B. username@domainname.com.

Installieren Sie für Scala das com.microsoft.aad.adal4j Artefakt.

Installieren Sie für Python die adal Bibliothek. Diese Bibliothek ist über pip verfügbar.

Beispiele finden Sie in den Beispielnotizbüchern.