Spark 連接器是 SQL 資料庫的一個高效能函式庫,讓你能從 SQL Server、Azure SQL 資料庫以及 Fabric 中的 SQL 資料庫讀取與寫入。 連接器提供下列功能:
- 使用 Spark 在 Azure SQL Database、Azure SQL 受控執行個體、Azure VM 上的 SQL Server,以及 Fabric 中的 SQL 資料庫上執行大型寫入與讀取操作。
- 當你使用表格或檢視時,連接器支援設定在 SQL 引擎層級的安全模型。 這些模型包括物件層級安全性(OLS)、資料列層級安全性(RLS),以及資料行層級安全性(CLS)。
連接器是預先安裝在 Fabric 執行環境裡,所以不需要另外安裝。
Authentication
Microsoft Entra 認證已整合於 Fabric 中。
- 當你登入 Fabric 工作區時,憑證會自動交給 SQL 引擎進行認證與授權。
- 需要在您的 SQL 資料庫引擎啟用並設定 Microsoft Entra ID。
- Microsoft Entra ID 已設定後,您的 Spark 程式碼不需要額外的設定。 憑證會自動映射。
你也可以使用 SQL 認證方法(透過指定 SQL 使用者名稱和密碼)或服務主體(提供 Azure 存取權杖用於應用程式認證)。
權限
要使用 Spark 連接器,你的身份——無論是使用者還是應用程式——都必須擁有目標 SQL 引擎所需的資料庫權限。 這些權限是從資料表和檢視讀取或寫入的必要條件。
對於 Azure SQL Database、Azure SQL 受控執行個體、Azure VM 上的 SQL Server:
- 執行操作的身份通常需要
db_datawriter和db_datareader權限,可選擇使用db_owner完全控制權限。
針對 Fabric 中的 SQL 資料庫:
- 身份通常需要
db_datawriter和db_datareader權限,此外可選擇db_owner權限。 - 該身分也需要至少具備在 Fabric 中於項目層級讀取 SQL 資料庫的權限。
備註
如果您使用服務主體帳戶,它可以在不涉及使用者上下文的情況下以應用程式形式運行,或者在啟用了使用者模擬後以使用者身份運行。 服務主體必須擁有你想執行的操作所需的資料庫權限。
使用與程式碼範例
本節提供程式碼範例,示範如何有效使用 Spark 連接器來管理 SQL 資料庫。 這些範例涵蓋多種情境,包括從 SQL 資料表讀取與寫入,以及配置連接器選項。
備註
在批量寫入前,所有進入的 Spark 資料必須符合目標 SQL 資料型態。 當你覆寫或建立資料表時,連接器會將 Spark TimestampType 和 TimestampNTZType 值映射成 SQL datetime2 ,而非 datetime。 Spark 時間戳記類型支援最多六位數的分秒精度,但 SQL datetime 支援三位數,這可能導致不匹配。
支援的選項
最小要求選項為 url as "jdbc:sqlserver://<server>:<port>;database=<database>;" 或 set spark.mssql.connector.default.url。
當
url提供時:- 永遠優先使用
url。 - 如果
spark.mssql.connector.default.url沒有設定,連接器會將其設定並在以後的用途中再次使用。
- 永遠優先使用
當未提供
url時:- 如果
spark.mssql.connector.default.url設定 ,連接器會使用火花設定中的值。 - 如果
spark.mssql.connector.default.url沒有設定,就會因為所需細節無法取得而產生錯誤。
- 如果
此連接器支援此處定義的選項: SQL DataSource JDBC 選項
連接器也支援以下選項:
| Option | 預設值 | Description |
|---|---|---|
reliabilityLevel |
「BEST_EFFORT」 | 控制插入操作的可靠性。 可能的值: BEST_EFFORT (預設、最快,執行器重新啟動時可能導致重複列)、 NO_DUPLICATES (較慢,確保即使執行器重新啟動也不會插入重複列)。 根據你對重複品的容忍度和效能需求來選擇。 |
isolationLevel |
「READ_COMMITTED」 | 設定 SQL 操作的交易隔離層級。 可能的值: READ_COMMITTED (預設值,防止讀取未提交的資料)、READ_UNCOMMITTED、REPEATABLE_READSNAPSHOTSERIALIZABLE。 較高的隔離等級可以減少並行性,但能提升資料一致性。 |
tableLock |
"錯誤" | 控制插入操作時是否使用 SQL Server TABLOCK 的表層級鎖定提示。 可能的數值: true (啟用 TABLOCK,可提升批量寫入效能)、 false (預設不使用 TABLOCK)。 設定為 true 可能會增加針對大型插入作業的吞吐量,但可能會減少對該表進行其他操作的並行性。 |
schemaCheckEnabled |
真實 | 控制 Spark DataFrame 與 SQL 資料表之間是否強制執行嚴格的結構驗證。 可能的值: true (預設值,強制嚴格的結構匹配), false (允許更多彈性,且可能跳過部分結構檢查)。 設定為 false 有助於解決結構不匹配,但若結構差異顯著,可能導致意想不到的結果。 |
其他 Bulk API 選項 可以設定為 Option DataFrame ,並在寫入時傳遞給 Bulk copy API。
書寫與讀取範例
以下程式碼使用自動 Microsoft Entra ID 驗證來示範這些操作:
- 寫入資料框架:
df.write.option("...", "...").mssql("<schema>.<table>")。 - 請閱讀一張表格:
spark.read.option("...", "...").mssql("<schema>.<table>")。 - 執行自訂查詢:
spark.read.option("...", "...").option("query", "<your-custom-query>").mssql()。
小提示
資料是為示範目的而內嵌產生的。 在生產環境中,你通常會從現有來源讀取資料,或建立更複雜的 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
你也可以選擇欄位、套用篩選器,並在讀取 SQL 資料庫引擎的資料時使用其他選項。
認證範例
以下範例展示了如何使用 Microsoft Entra ID 以外的認證方法,例如服務主體(存取權杖)和 SQL 認證。
備註
如前所述,Microsoft Entra ID 認證會在你登入 Fabric 工作區時自動處理,因此只有在情境需要時才需要使用這些方法。
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
支援的 DataFrame 儲存模式
當你從 Spark 寫入 SQL 資料庫資料時,可以選擇多種存檔模式。 儲存模式控制當目的資料表已存在時的資料如何寫入,並可能影響結構、資料與索引。 了解這些模式有助於避免意外的資料遺失或變更。
此連接器支援以下定義的選項: Spark Save 函式
ErrorIfExists (預設儲存模式):如果目標資料表存在,寫入會中止並回傳例外。 否則,會建立一個包含資料的新資料表。
忽略:如果目標資料表存在,寫入會忽略請求且不會回傳錯誤。 否則,會建立一個包含資料的新資料表。
覆寫:如果目標資料表存在,則丟棄該資料表,重新建立,並新增資料。
備註
當你使用
overwrite時,會失去原始的資料表結構(尤其是 MSSQL 獨有的資料型別)和資料表索引。 這個結構會被從你的 Spark DataFrame 推斷出的結構所取代。 為避免失去結構與索引,請加入.option("truncate", true)。附加:若目標資料表存在,則會附加新資料。 否則,會建立一個包含資料的新資料表。
疑難排解
當程序結束時,你的 Spark 讀取操作的輸出會出現在該儲存格的輸出區域。 從 com.microsoft.sqlserver.jdbc.SQLServerException 出現的錯誤直接來自 SQL Server。 您可以在 Spark 應用程式日誌中找到詳細的錯誤資訊。
批次寫入需要輸入資料來符合目標 SQL 資料型態。 如果資料不符,你可能會收到以下錯誤:
Caused by: com.microsoft.sqlserver.jdbc.SQLServerException: The service has encountered an error processing your request. Please try again. Error code 4815.
例如,當目標 SQL 資料表使用 datetime,但輸入的 Spark TimestampType 值(如 2025-01-01 10:30:00.123456)的精度高於 datetime 支援值時,可能會發生此錯誤。
要解決這個錯誤,請採用以下其中一種方法:
- 將輸入資料鑄造、截斷或轉換,使其符合目標 SQL 資料型態。 例如,將數值截斷為三位數的分秒精度:
2025-01-01 10:30:00.123000。 - 允許連接器透過設定
.option("truncate", false)來重建資料表結構。 連接器將 Spark 時間戳記型別對應至 SQLdatetime2。