Konektor Spark pro databáze SQL

Spark konektor pro SQL databáze je vysoce výkonná knihovna, která umožňuje číst a zapisovat do SQL Server, Azure SQL databází a SQL databází ve Fabric. Konektor nabízí následující funkce:

  • Používejte Spark k provádění velkých operací zápisu a čtení na Azure SQL Database, Azure SQL Managed Instance, SQL Server na Azure VM a SQL databázích ve Fabric.
  • Pokud používáte tabulku nebo zobrazení, konektor podporuje modely zabezpečení nastavené na úrovni modulu SQL. Mezi tyto modely patří zabezpečení na úrovni objektů (OLS), zabezpečení na úrovni řádků (RLS) a zabezpečení na úrovni sloupců (CLS).

Konektor je předinstalovaný v modulu runtime Fabric, takže ho nemusíte instalovat samostatně.

Autentizace

Autentizace Microsoft Entra je integrována s Fabric.

  • Když se přihlásíte k pracovnímu prostoru Fabric, vaše přihlašovací údaje se automaticky předají modulu SQL za účelem ověřování a autorizace.
  • Vyžaduje, aby bylo v databázovém stroji SQL povolené a nakonfigurované ID Microsoft Entra.
  • Pokud je v kódu Sparku nastaveno ID Microsoft Entra, není potřeba žádná další konfigurace. Přihlašovací údaje se automaticky mapují.

Můžete také použít metodu ověřování SQL (zadáním uživatelského jména a hesla SQL) nebo instančního objektu (poskytnutím přístupového tokenu Azure pro ověřování na základě aplikace).

Povolení

Pokud chcete používat konektor Spark, musí mít vaše identita (ať už uživatel nebo aplikace) potřebná oprávnění k databázi pro cílový modul SQL. Tato oprávnění jsou nutná ke čtení z tabulek a zobrazení a k zápisu do tabulek a zobrazení.

Pro Azure SQL Database, Azure SQL Managed Instance a SQL Server na virtuálním počítači Azure:

  • Identita, na které se operace spouští, obvykle potřebuje oprávnění db_datawriter a db_datareader, a volitelně db_owner pro úplné řízení.

Pro databázi SQL v prostředí Fabric:

  • Identita obvykle potřebuje db_datawriter a db_datareader oprávnění a volitelně db_owner.
  • Identita také vyžaduje alespoň oprávnění ke čtení v SQL databázi ve Fabric na úrovni položek.

Poznámka:

Pokud používáte hlavní službu, může se spustit jako aplikace (bez kontextu uživatele) nebo jako uživatel, pokud je povoleno zosobnění uživatele. Spouštěcí služba musí mít požadovaná oprávnění k databázi pro operace, které chcete provést.

Příklady použití a kódu

V této části uvádíme příklady kódu, které ukazují, jak efektivně používat konektor Sparku pro databáze SQL. Tyto příklady pokrývají různé scénáře, včetně čtení z tabulek SQL a zápisu do tabulek SQL a konfigurace možností konektoru.

Poznámka:

Před hromadným zápisem musí všechna příchozí data ze Sparku odpovídat cílovým SQL datovým typům. Když přepíšete nebo vytvoříte tabulku, konektor mapuje hodnoty Spark TimestampType a TimestampNTZType na hodnoty SQL datetime2 místo datetime. Typy časových razítek ve Sparku podporují až šest číslic zlomku sekundy, ale SQL datetime podporuje tři, což může způsobit neshodu.

Podporované možnosti

Minimální požadovaná možnost je url jako "jdbc:sqlserver://<server>:<port>;database=<database>;" nebo spark.mssql.connector.default.url.

  • Pokud je url k dispozici:

    • Vždy používejte url jako první předvolbu.
    • Pokud spark.mssql.connector.default.url není nastavený, konektor ho nastaví a znovu použije pro budoucí použití.
  • Pokud url není k dispozici:

    • Pokud spark.mssql.connector.default.url je nastavená, konektor použije hodnotu z konfigurace Sparku.
    • Pokud spark.mssql.connector.default.url není nastavená, vyvolá se chyba, protože požadované podrobnosti nejsou k dispozici.

Tento konektor podporuje možnosti definované tady: Možnosti SQL DataSource JDBC

Konektor také podporuje následující možnosti:

Možnost Výchozí hodnota Description
reliabilityLevel Nejlepší Úsilí Řídí spolehlivost operací vkládání. Možné hodnoty: BEST_EFFORT (výchozí, nejrychlejší, může vést k duplicitním řádkům, pokud se exekutor restartuje), NO_DUPLICATES (pomalejší, zajistí se, že se nebudou vkládat žádné duplicitní řádky, i když se exekutor restartuje). Zvolte si na základě vaší tolerance duplicitních položek a požadavků na výkon.
isolationLevel "READ_COMMITTED" Nastaví úroveň izolace transakcí pro operace SQL. Možné hodnoty: READ_COMMITTED (výchozí hodnota brání čtení nepotvrzených dat), READ_UNCOMMITTED, REPEATABLE_READ, SNAPSHOT, . SERIALIZABLE Vyšší úrovně izolace můžou snížit souběžnost, ale zlepšit konzistenci dat.
tableLock nepravda Určuje, zda se během operací vložení používá nápověda k uzamčení na úrovni tabulky SQL Server TABLOCK. Možné hodnoty: true (povolí TABLOCK, což může zlepšit výkon hromadného zápisu), false (výchozí hodnota nepoužívá TABLOCK). Nastavení na true může zvýšit propustnost při velkých vkládáních, ale také snížit souběžnost pro jiné operace na tabulce.
schemaCheckEnabled pravda Určuje, jestli se mezi Sparkem DataFrame a tabulkou SQL vynucuje přísné ověřování schématu. Možné hodnoty: true (výchozí nastavení vynucuje striktní porovnávání schémat), false (umožňuje větší flexibilitu a může přeskočit některé kontroly schématu). Nastavení na false může pomoci při neshodách schémat, ale pokud se struktury podstatně liší, může vést k neočekávaným výsledkům.

Další možnosti hromadného rozhraní API lze nastavit jako parametry na DataFrame a předat je API hromadného kopírování během zápisu.

Příklad zápisu a čtení

Následující kód používá automatickou autentizaci Microsoft Entra ID k demonstraci těchto operací:

  • Zapište datový rámec: df.write.option("...", "...").mssql("<schema>.<table>").
  • Přečtěte si tabulku: spark.read.option("...", "...").mssql("<schema>.<table>").
  • Spusť vlastní dotaz spark.read.option("...", "...").option("query", "<your-custom-query>").mssql(): .

Návod

Data jsou vytvořena inline pro demonstrační účely. V produkčním scénáři byste obvykle četli data z existujícího zdroje nebo vytvořili složitější 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

Můžete také vybrat sloupce, použít filtry a použít další možnosti při čtení dat z databázového stroje SQL.

Příklady ověřování

Následující příklady ukazují, jak používat jiné metody ověřování než Microsoft Entra ID, jako je instanční objekt (přístupový token) a ověřování SQL.

Poznámka:

Jak už bylo zmíněno dříve, ověřování Microsoft Entra ID se zpracovává automaticky při přihlášení k pracovnímu prostoru Fabric, takže tyto metody stačí použít jenom v případě, že je váš scénář vyžaduje.

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

Podporované režimy ukládání datového rámce

Při zápisu dat ze Sparku do databází SQL si můžete vybrat z několika režimů ukládání. Režimy ukládání určují, jak se data zapisují, když už cílová tabulka existuje, a může ovlivnit schéma, data a indexování. Pochopení těchto režimů vám pomůže vyhnout se neočekávané ztrátě nebo změnám dat.

Tento konektor podporuje možnosti definované zde: Funkce Spark Save

  • ErrorIfExists (výchozí režim ukládání): Pokud cílová tabulka existuje, zápis se přeruší a vrátí se výjimka. V opačném případě se vytvoří nová tabulka s daty.

  • Ignorovat: Pokud cílová tabulka existuje, zápis požadavek ignoruje a nevrátí chybu. V opačném případě se vytvoří nová tabulka s daty.

  • Přepsání: Pokud cílová tabulka existuje, tabulka se zahodí, znovu vytvoří a připojí se nová data.

    Poznámka:

    Když použijete overwrite, ztratíte původní tabulkové schéma (zejména datové typy exkluzivní pro MSSQL) a indexy tabulek. Schéma je nahrazeno schématem odvozeným z vašeho Spark DataFrame. Aby nedošlo ke ztrátě schématu a indexů, přidejte .option("truncate", true).

  • Připojení: Pokud cílová tabulka existuje, připojí se k ní nová data. V opačném případě se vytvoří nová tabulka s daty.

Troubleshoot

Po dokončení procesu se výstup operace čtení Sparku zobrazí ve výstupní oblasti buňky. Chyby z com.microsoft.sqlserver.jdbc.SQLServerException pocházejí přímo ze SQL Serveru. Podrobné informace o chybách najdete v protokolech aplikace Spark.

Hromadné zápisy vyžadují, aby příchozí data odpovídala cílovým SQL datovým typům. Pokud data nesedí, můžete dostat tuto chybu:

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

Například k této chybě může dojít, když cílová tabulka SQL používá typ datetime, ale příchozí hodnota Spark typu TimestampType, například 2025-01-01 10:30:00.123456, má vyšší přesnost, než jakou podporuje datetime.

K odstranění chyby použijte jeden z těchto přístupů:

  • Přenášejte, zkracujte nebo transformujte příchozí data tak, aby odpovídala cílovému typu SQL dat. Například zkraťte hodnotu na tři číslice přesnosti zlomků sekundy: 2025-01-01 10:30:00.123000.
  • Nechte konektor znovu vytvořit schéma tabulky nastavením .option("truncate", false). Konektor mapuje časový typ Sparku na SQL datetime2.