Il connettore Spark per i database SQL

Il connettore Spark per database SQL è una libreria ad alte prestazioni che ti permette di leggere e scrivere su database SQL Server, Azure SQL e database SQL in Fabric. Il connettore offre le funzionalità seguenti:

  • Usa Spark per eseguire grandi operazioni di scrittura e lettura su database SQL di Azure, Istanza gestita di SQL di Azure, SQL Server su Azure VM e database SQL in Fabric.
  • Quando si usa una tabella o una vista, il connettore supporta i modelli di sicurezza impostati a livello di motore SQL. Questi modelli includono la protezione a livello di oggetto (OLS), la protezione a livello di riga (RLS) e la protezione a livello di colonna (CLS).

Il connettore è preinstallato nel runtime di Fabric, quindi non è necessario installarlo separatamente.

Authentication

L'autenticazione Microsoft Entra è integrata con Fabric.

  • Quando si accede all'area di lavoro Infrastruttura, le credenziali vengono passate automaticamente al motore SQL per l'autenticazione e l'autorizzazione.
  • Richiede che Microsoft Entra ID sia abilitato e configurato nel motore di database SQL.
  • Non è necessaria alcuna configurazione aggiuntiva nel codice Spark se è configurato Microsoft Entra ID. Le credenziali vengono mappate automaticamente.

È anche possibile usare il metodo di autenticazione SQL (specificando un nome utente e una password SQL) o un'entità servizio (fornendo un token di accesso di Azure per l'autenticazione basata su app).

Permissions

Per usare il connettore Spark, la tua identità, sia essa utente o applicazione, deve disporre delle autorizzazioni di database necessarie per il motore SQL di destinazione. Queste autorizzazioni sono necessarie per leggere o scrivere su tabelle e viste.

Per il database SQL di Azure, Istanza gestita di SQL di Azure e SQL Server nella macchina virtuale di Azure:

  • L'identità che esegue l'operazione in genere necessita delle autorizzazioni db_datawriter e db_datareader, e facoltativamente db_owner per il controllo completo.

Per un database SQL in Fabric:

  • L'identità in genere richiede le autorizzazioni db_datawriter e db_datareader, e facoltativamente db_owner.
  • L'identità richiede anche almeno il permesso di lettura sul database SQL in Fabric a livello di elemento.

Annotazioni

Se si utilizza un principale del servizio, può operare come app (senza contesto utente) o come utente nel caso in cui sia abilitata la rappresentazione dell'utente. Il principal del servizio deve disporre delle autorizzazioni necessarie sul database per le operazioni da eseguire.

Esempi di utilizzo e codice

In questa sezione vengono forniti esempi di codice per illustrare come usare il connettore Spark per i database SQL in modo efficace. Questi esempi illustrano vari scenari, tra cui la lettura e la scrittura in tabelle SQL e la configurazione delle opzioni del connettore.

Annotazioni

Prima di una scrittura in massa, tutti i dati Spark in arrivo devono adattarsi ai tipi di dati SQL target. Quando sovrascrivi o crei una tabella, il connettore mappa Spark TimestampType e TimestampNTZType i valori in SQL datetime2 invece che datetime. I tipi di timestamp Spark supportano fino a sei cifre con precisione di frazione di secondo, ma SQL datetime ne supporta tre, il che può causare un disadattamento.

Opzioni supportate

L'opzione minima richiesta è url come "jdbc:sqlserver://<server>:<port>;database=<database>;" o spark.mssql.connector.default.url.

  • Quando viene fornito url:

    • Usare url sempre come prima preferenza.
    • Se spark.mssql.connector.default.url non è impostato, il connettore lo imposta e lo riutilizza per un utilizzo futuro.
  • Quando url non viene specificato:

    • Se spark.mssql.connector.default.url è impostato, il connettore usa il valore della configurazione spark.
    • Se spark.mssql.connector.default.url non è impostato, viene generato un errore perché i dettagli richiesti non sono disponibili.

Questo connettore supporta le opzioni definite qui: Opzioni JDBC di SQL DataSource

Il connettore supporta anche le opzioni seguenti:

Opzione Valore predefinito Description
reliabilityLevel Miglior_Sforzo Controlla l'affidabilità delle operazioni di inserimento. Valori possibili: BEST_EFFORT (impostazione predefinita, più veloce, potrebbe comportare righe duplicate se un executor viene riavviato), NO_DUPLICATES (più lento, garantisce che non vengano inserite righe duplicate anche se un executor viene riavviato). Scegliere in base alla tolleranza per i duplicati e le esigenze di prestazioni.
isolationLevel "READ_COMMITTED" Imposta il livello di isolamento delle transazioni per le operazioni SQL. Valori possibili: READ_COMMITTED (impostazione predefinita, impedisce la lettura dei dati non inviati), READ_UNCOMMITTED, REPEATABLE_READSNAPSHOT, , SERIALIZABLE. Livelli di isolamento più elevati possono ridurre la concorrenza, ma migliorare la coerenza dei dati.
tableLock falso Controlla se l'hint di blocco a livello di tabella TABLOCK di SQL Server viene utilizzato durante le operazioni di inserimento. Valori possibili: true (abilita TABLOCK, che può migliorare le prestazioni di scrittura bulk), false (impostazione predefinita, non usa TABLOCK). L'impostazione su true potrebbe aumentare la velocità effettiva per inserimenti di grandi dimensioni, ma può ridurre la possibilità di eseguire operazioni concorrenti sulla tabella.
schemaCheckEnabled vero Controlla se la convalida dello schema rigorosa viene applicata tra Spark DataFrame e la tabella SQL. Valori possibili: true (impostazione predefinita, applica la corrispondenza dello schema rigorosa), false (consente una maggiore flessibilità e potrebbe ignorare alcuni controlli dello schema). Impostare su false può essere utile per le discrepanze dello schema, ma potrebbe causare risultati imprevisti se le strutture differiscono in modo significativo.

Altre opzioni dell'API bulk possono essere impostate come opzioni in DataFrame e vengono passate alle API di copia bulk in scrittura.

Esempio di scrittura e lettura

Il seguente codice utilizza l'autenticazione automatica Microsoft Entra ID per dimostrare queste operazioni:

  • Scrivi un DataFrame: df.write.option("...", "...").mssql("<schema>.<table>").
  • Leggi una tabella: spark.read.option("...", "...").mssql("<schema>.<table>").
  • Esegui una query personalizzata: spark.read.option("...", "...").option("query", "<your-custom-query>").mssql().

Suggerimento

I dati sono creati in linea a scopo dimostrativo o illustrativo. In uno scenario di produzione, in genere si leggerebbero i dati da un'origine esistente o si creerebbe un oggetto più complesso DataFrame.

import com.microsoft.sqlserver.jdbc.spark
url = "jdbc:sqlserver://<server>:<port>;database=<database>;"
row_data = [("Alice", 1),("Bob", 2),("Charlie", 3)]
column_header = ["Name", "Age"]
df = spark.createDataFrame(row_data, column_header)
df.write.mode("overwrite").option("url", url).mssql("dbo.publicExample")
spark.read.option("url", url).mssql("dbo.publicExample").show()
spark.read.option("url", url).option("query", "SELECT * FROM dbo.publicExample WHERE Age = 3").mssql().show() # Read with a custom query

url = "jdbc:sqlserver://<server>:<port>;database=<database2>;" # different database
df.write.mode("overwrite").option("url", url).mssql("dbo.tableInDatabase2") # default url is updated
spark.read.mssql("dbo.tableInDatabase2").show() # no url option specified and will use database2

È anche possibile selezionare colonne, applicare filtri e usare altre opzioni quando si leggono i dati dal motore di database SQL.

Esempi di autenticazione

Gli esempi seguenti illustrano come usare metodi di autenticazione diversi da Microsoft Entra ID, ad esempio l'entità servizio (token di accesso) e l'autenticazione SQL.

Annotazioni

Come accennato in precedenza, l'autenticazione di Microsoft Entra ID viene gestita automaticamente quando si accede all'area di lavoro di Fabric, quindi è necessario usare questi metodi solo se lo scenario li richiede.

import com.microsoft.sqlserver.jdbc.spark
url = "jdbc:sqlserver://<server>:<port>;database=<database>;"
row_data = [("Alice", 1),("Bob", 2),("Charlie", 3)]
column_header = ["Name", "Age"]
df = spark.createDataFrame(row_data, column_header)

from azure.identity import ClientSecretCredential
credential = ClientSecretCredential(tenant_id="", client_id="", client_secret="") # service principal app
scope = "https://database.windows.net/.default"
token = credential.get_token(scope).token

df.write.mode("overwrite").option("url", url).option("accesstoken", token).mssql("dbo.publicExample")
spark.read.option("accesstoken", token).mssql("dbo.publicExample").show()
spark.read.option("accesstoken", token).option("query", "SELECT * FROM dbo.publicExample WHERE Age = 3").mssql().show() # Read with a custom query

Modalità supportate di salvataggio del DataFrame

Quando si scrivono dati da Spark a database SQL, è possibile scegliere tra diverse modalità di salvataggio. Le modalità di salvataggio controllano la modalità di scrittura dei dati quando la tabella di destinazione esiste già e possono influire su schema, dati e indicizzazione. La comprensione di queste modalità consente di evitare la perdita o le modifiche impreviste dei dati.

Questo connettore supporta le opzioni definite qui: Funzioni di salvataggio Spark

  • ErrorIfExists (modalità di salvataggio predefinita): se la tabella di destinazione esiste, la scrittura viene interrotta e viene restituita un'eccezione. In caso contrario, viene creata una nuova tabella con i dati.

  • Ignora: se la tabella di destinazione esiste, la scrittura ignora la richiesta e non restituisce un errore. In caso contrario, viene creata una nuova tabella con i dati.

  • Sovrascrivi: se la tabella di destinazione esiste, la tabella viene eliminata, creata nuovamente e i nuovi dati vengono aggiunti.

    Annotazioni

    Quando usi overwrite, perdi lo schema originale delle tabelle (specialmente i tipi di dati esclusivi MSSQL) e gli indici delle tabelle. Lo schema viene sostituito dallo schema dedotto dal tuo DataFrame Spark. Per evitare di perdere lo schema e gli indici, aggiungi .option("truncate", true).

  • Accoda: se la tabella di destinazione esiste, vengono aggiunti nuovi dati. In caso contrario, viene creata una nuova tabella con i dati.

Troubleshoot

Al termine del processo, l'output dell'operazione di lettura Spark sarà visibile nell'area di output della cella. com.microsoft.sqlserver.jdbc.SQLServerException Gli errori provengono direttamente da SQL Server. È possibile trovare informazioni dettagliate sull'errore nei log dell'applicazione Spark.

Le scritture in massa richiedono che i dati in arrivo si adattino ai tipi di dati SQL di destinazione. Se i dati non corrispondono, potresti ricevere questo errore:

Caused by: com.microsoft.sqlserver.jdbc.SQLServerException: The service has encountered an error processing your request. Please try again. Error code 4815.

Ad esempio, questo errore può verificarsi quando la tabella SQL di destinazione utilizza datetime, ma un valore TimestampType di Spark in ingresso, come 2025-01-01 10:30:00.123456, ha una precisione superiore a quella supportata da datetime.

Per risolvere l'errore, si utilizza uno di questi approcci:

  • Converti, tronca o trasforma i dati in arrivo per adattarli al tipo di dati SQL di destinazione. Ad esempio, tronca il valore a tre cifre per la precisione delle frazioni di secondo: 2025-01-01 10:30:00.123000.
  • Permette al connettore di ricreare lo schema della tabella impostando .option("truncate", false). Il connettore mappa il tipo di timestamp di Spark a SQL datetime2.