Microsoft JDBC 驅動程式用於 Microsoft Fabric 資料工程

JDBC(Java 資料庫連接性)是一項廣泛採用的標準,使客戶端應用程式能夠連接並操作來自資料庫及大數據平台的資料。

Microsoft JDBC Fabric 資料工程驅動程式讓您能以 JDBC 標準的可靠性與簡潔性,在 Fabric 中連接、查詢並管理 Spark 工作負載。 該驅動程式建立在 Fabric 的 Livy API 之上,為您的 Java 應用程式與 BI 工具提供安全且靈活的 Spark SQL 連接。 此整合讓您能直接提交並執行 Spark 程式碼,無需另建筆記本或 Spark 工作定義項目。 此驅動程式相容於熱門的 JDBC 用戶端,如 DbVisualizer 和 DBeaver,以及支援 JDBC 連接的 BI 工具,包括 Tableau。

主要功能

  • JDBC 4.2 API:實作 JDBC 4.2 API 以支援 Spark SQL 連線,但受 Spark 與 Livy 的限制,例如不支援 JDBC 交易,以及僅向前、唯讀的結果集
  • Microsoft Entra ID 認證:多重認證流程,包括互動式、用戶端憑證及憑證式認證
  • 高並行性會話:為並行工作負載提供自願加入的 Fabric 會話共享
  • 顯式連線集區:選擇 HikariCP(標準的集區管理器)或 LivyBuiltInPooledDataSource;直接進入點永遠不會使用連線集區
  • Spark SQL 原生查詢支援:直接執行 Spark SQL 語句而無需轉換
  • 全面的資料型態支援:支援所有 Spark SQL 資料型態,包括複雜型別(ARRAY、MAP、STRUCT)
  • 非同步結果集預取:背景資料載入以提升效能
  • 斷路器模式:防止連鎖故障,有自動重試功能以保護系統
  • 自動重新連線:連線失敗時透明的會話恢復
  • 進階重試邏輯:透過指數退縮與會話恢復來提升韌性
  • 代理支援:企業環境的 HTTP 與 SOCKS 代理設定

先決條件

在使用 Microsoft JDBC Fabric 資料工程驅動程式前,請確保您具備:

  • Java 開發套件(JDK):Java 11、17 和 21(推薦使用 Java 21)
  • Fabric 存取權:Fabric 工作區的存取權
  • Microsoft Entra ID 憑證:適合的憑證用於認證
  • 工作區與湖屋 ID:用於 Fabric 工作區與湖屋的 GUID 識別碼

下載與安裝

Microsoft JDBC 驅動程式適用於 Fabric 資料工程 2.0.1 版本支援 Java 11、17 及 21。 使用最新版本。

  1. 請從上方連結下載壓縮檔或 tar 檔。
  2. 解壓下載的檔案以存取驅動程式的 JAR 檔案。
  3. 選擇與你 JRE 版本相符的 JAR 檔案:
    • 針對 Java 11: ms-sparksql-jdbc-2.0.1.jre11.jar
    • 針對 Java 17: ms-sparksql-jdbc-2.0.1.jre17.jar
    • 針對 Java 21: ms-sparksql-jdbc-2.0.1.jre21.jar
  4. 將所選的 JAR 檔案加入你的應用程式類別路徑。
  5. 對於 JDBC 用戶端,請設定 JDBC 驅動類別: com.microsoft.spark.livy.jdbc.LivyDriver

快速入門範例

此範例示範如何使用 Microsoft JDBC 驅動程式連接 Fabric 並執行查詢Fabric資料工程。 執行此程式碼前,請確保你完成前置條件並安裝驅動程式。 藉由使用 AuthFlow=2,驅動程式會使用 DefaultAzureCredential。 要使用 Azure CLI 作為憑證來源,安裝 Azure CLI 並執行 az login.

import java.sql.*;

public class QuickStartExample {
    public static void main(String[] args) {
        // Connection string with required parameters
        String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
                     "FabricWorkspaceID=<workspace-id>;" +
                     "FabricLakehouseID=<lakehouse-id>;" +
                     "AuthFlow=2;" +  // DefaultAzureCredential; Azure CLI can be one source
                     "LogLevel=INFO";
        
        try (Connection conn = DriverManager.getConnection(url)) {
            // Execute a simple query
            try (Statement stmt = conn.createStatement();
                 ResultSet rs = stmt.executeQuery("SELECT 'Hello from Fabric!' as message")) {
                
                if (rs.next()) {
                    System.out.println(rs.getString("message"));
                }
            }
        } catch (SQLException e) {
            e.printStackTrace();
        }
    }
}

連接字串格式

基本連接串

Microsoft 的 Fabric 資料工程 JDBC 驅動程式使用以下 連接字串 格式:

jdbc:fabricspark://<hostname>[:<port>][;<parameter1>=<value1>;<parameter2>=<value2>;...]

連接字串組件

元件 Description Example
通訊協定 JDBC URL 協定識別碼 jdbc:fabricspark://
主機名稱 Fabric 端點主機名稱 api.fabric.microsoft.com
通訊埠 可選埠號(預設:443) :443
參數 分號分隔鍵=值對 FabricWorkspaceID=<guid>

連接字串範例

基本連線(互動式瀏覽器驗證)

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;AuthFlow=1

Spark 資源設定

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;DriverCores=4;DriverMemory=4g;ExecutorCores=4;ExecutorMemory=8g;NumExecutors=2;AuthFlow=2

與 Spark Session Properties 合作

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;spark.sql.adaptive.enabled=true;spark.sql.shuffle.partitions=200;AuthFlow=2

採用高併發模式

在驅動程式版本 1.1.0 及以上版本中,Microsoft Fabric 連線可使用高並行(HC)模式。 加入 hcEnabled=true 以取得 HC 工作階段,而非傳統的 Livy 工作階段:

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;AuthFlow=2;hcEnabled=true

你可以選擇用會話標籤來識別工作負載,並選擇 Fabric 環境:

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;AuthFlow=2;hcEnabled=true;sessionTag=orders_etl;FabricEnvironmentID=<environment-id>

HC 模式需主動啟用。 若 hcEnabled 省略或設為 false,驅動程式則使用經典的 Livy 會話路徑。

完整配置指引請參閱 High-Concurrency(HC)模式 章節。


Authentication

Microsoft JDBC Fabric Data Engineering 驅動程式支援透過 Microsoft Entra ID(前稱 Azure Active Directory)進行多種認證方法。 認證是透過 AuthFlow 連接字串中的參數來設定的。

認證流程

身份驗證流程 驗證方法 用例
1 互動式瀏覽器 使用 OAuth 2.0 的互動式使用者驗證
2 預設 Azure 認證鏈 開發與管理應用程式認證;Azure CLI 可以是單一憑證來源
3 用戶端秘密憑證(服務主體) 自動化/服務對服務認證
4 用戶端憑證 基於憑證的服務主體認證
5 存取令牌 預先取得的承載者存取憑證

互動式瀏覽器認證

最佳用途: 開發與互動應用

String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=1;" +
             "LogLevel=INFO";

Connection conn = DriverManager.getConnection(url);

參數:

  • AuthFlow=1: 規範互動式瀏覽器認證

態度:

  • 開啟瀏覽器視窗以進行使用者驗證
  • 會暫存憑證以供後續連線使用,直到連線過期
  • 適合單用戶應用

用戶端憑證或服務主體認證

最佳用途: 自動化服務與背景工作

String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=3;" +  
             "AuthClientID=<client-id>;" +
             "AuthClientSecret=<client-secret>;" +
             "AuthTenantID=<tenant-id>;" +
             "LogLevel=INFO";

Connection conn = DriverManager.getConnection(url);

所需參數:

  • AuthFlow=3: 規定用戶端憑證認證
  • AuthClientID:來自 Microsoft Entra ID 的應用程式(用戶端)ID
  • AuthClientSecret:Microsoft Entra ID 的用戶端密碼
  • AuthTenantID:Microsoft Entra租戶識別碼

最佳實踐:

  • 安全儲存秘密(Azure Key Vault, environment variables)
  • 盡可能使用受管理身份
  • 要定期輪換密鑰

基於憑證的認證

最佳應用: 需要憑證式認證的企業應用

String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=4;" +  
             "AuthClientID=<client-id>;" +
             "AuthCertificatePath=/path/to/certificate.pem;" +
             "AuthCertificatePassword=<certificate-password>;" +
             "AuthTenantID=<tenant-id>;" +
             "LogLevel=INFO";

Connection conn = DriverManager.getConnection(url);

所需參數:

  • AuthFlow=4:規定基於憑證的認證
  • AuthClientID: 應用程式(用戶端)ID
  • AuthCertificatePath: PEM 憑證檔案的路徑
  • AuthCertificatePassword: 證書密碼
  • AuthTenantID:Microsoft Entra租戶識別碼

存取令牌認證

最佳用途: 自訂認證場景

// Acquire token through custom mechanism
String accessToken = acquireTokenFromCustomSource();

String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=5;" +  // Access token authentication
             "AuthAccessToken=" + accessToken + ";" +
             "LogLevel=INFO";

Connection conn = DriverManager.getConnection(url);

認證快取

驅動程式會快取認證憑證,並根據憑證到期時間刷新。 已剖析的 AuthEnableCaching 和 AuthCacheTTLMS 屬性在 2.0.0 版中無法控制此行為。

配置參數

必要參數

這些參數必須存在於每個連接字串中:

參數 類型 Description Example
FabricWorkspaceID 通用唯一識別碼 (UUID) Fabric 工作區識別碼 <workspace-id>
FabricLakehouseID 通用唯一識別碼 (UUID) Fabric 湖倉識別碼 <lakehouse-id>
AuthFlow 整數 認證流程類型(1-5) 2

選擇性參數

API 版本設定

參數 類型 預設 Description
FabricVersion 繩子 v1 Fabric API 版本
LivyApiVersion 繩子 2023-12-01 Livy API 版本

環境設定

參數 類型 預設 Description
FabricEnvironmentID 通用唯一識別碼 (UUID) None 用於 Spark 會話引用環境項目的 Fabric 環境識別碼

高並行會話配置

設定 hcEnabled=true 為使用高效能模式。 剩餘的 HC 屬性僅在啟用 HC 模式時使用。

參數 類型 預設 Description
hcEnabled 布林值 false 使用 Fabric 高並行會話。
sessionTag 繩子 None 可選的伺服器端會話標籤。 請使用1至64的字母、數字、底線或連字號。
hcAcquireTimeoutSeconds 整數 300 等待HC會議準備好的最長時間。 有效時間範圍為120至3,600秒。
hcAcquirePollingIntervalMs 整數 1000 HC 會話狀態檢查之間的間隔。 接受的數值範圍為 50 到 30,000 毫秒,執行時會被壓縮到 100 到 5,000 毫秒。

Example:

hcEnabled=true;sessionTag=interactive_reporting;hcAcquireTimeoutSeconds=600;hcAcquirePollingIntervalMs=2000;FabricEnvironmentID=<environment-id>
  • FabricEnvironmentID 適用於硬核模式和經典場次
  • EnvironmentID 也可作為跨驅動程式別名。 屬性名稱不區分大小寫。

Important

僅將 sessionTag 用於非敏感的操作標籤。 不要包含機密、存取權杖、個人資料、客戶識別碼或查詢文字。

完整配置指引請參閱 High-Concurrency(HC)模式 章節。

Spark 配置

會話資源配置

配置 Spark 會話資源以達到最佳效能:

參數 類型 預設 Description Example
DriverCores 整數 Spark 預設 驅動程式的 CPU 核心數量 4
DriverMemory 繩子 Spark 預設 驅動程式的記憶體配置 4g
ExecutorCores 整數 Spark 預設 每個執行器擁有的 CPU 核心數 4
ExecutorMemory 繩子 Spark 預設 每個執行者的記憶體配置 8g
NumExecutors 整數 Spark 預設 執行程式數目 2
SessionName 繩子 自動生成 自訂會話名稱 MySparkSession

Example:

DriverCores=4;DriverMemory=4g;ExecutorCores=4;ExecutorMemory=8g;NumExecutors=2; SessionName=MySparkSession

自訂 Spark Session 屬性

任何帶有前綴 spark. 的參數都會自動套用到 Spark 會話:

Spark 配置範例:

spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.shuffle.partitions=200
spark.sql.autoBroadcastJoinThreshold=10485760
spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.minExecutors=1
spark.dynamicAllocation.maxExecutors=10
spark.executor.memoryOverhead=1g

原生執行引擎(NEE):

spark.nee.enabled=true

完整範例:

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<guid>;FabricLakehouseID=<guid>;DriverMemory=4g;ExecutorMemory=8g;NumExecutors=2;spark.sql.adaptive.enabled=true;spark.nee.enabled=true;AuthFlow=2

HTTP 用戶端連線設定

配置驅動程式的 HTTP 傳輸連線以達到最佳網路效能。 這些設定無法設定或管理 JDBC 連線池:

參數 類型 預設 Description
HttpMaxTotalConnections 整數 100 最大總 HTTP 連線數
HttpMaxConnectionsPerRoute 整數 20 每條路線的最大連接數
HttpConnectionTimeoutInSeconds 整數 30 連線逾時
HttpSocketTimeoutInSeconds 整數 60 插槽讀取逾時
HttpReadTimeoutInSeconds 整數 60 HTTP 讀取逾時
HttpConnectionRequestTimeoutSeconds 整數 10 來自池中的連線請求逾時
HttpEnableKeepAlive 布林值 true 啟用 HTTP 保持生命
HttpKeepAliveTimeoutSeconds 整數 60 維持連線逾時
HttpFollowRedirects 布林值 true 請遵循 HTTP 重定向
HttpUseAsyncIO 布林值 true 使用非同步 HTTP I/O

Example:

HttpMaxTotalConnections=200;HttpMaxConnectionsPerRoute=100;HttpConnectionTimeoutInSeconds=60

代理設定

為企業環境設定 HTTP 與 SOCKS 代理設定:

參數 類型 預設 Description
UseProxy 布林值 假的 啟用代理
ProxyTransport 繩子 http 代理傳輸類型(http/tcp)
ProxyHost 繩子 None 代理主機名稱
ProxyPort 整數 None 代理埠
ProxyAuthEnabled 布林值 假的 啟用代理驗證
ProxyUsername 繩子 None 代理驗證使用者名稱
ProxyPassword 繩子 None 代理驗證密碼
ProxyAuthScheme 繩子 basic 身份驗證機制(basic/digest/ntlm)
ProxySocksVersion 整數 5 SOCKS 版本(4/5)

HTTP 代理範例:

UseProxy=true;ProxyTransport=http;ProxyHost=proxy.company.com;ProxyPort=8080;ProxyAuthEnabled=true;ProxyUsername=user;ProxyPassword=pass

SOCKS 代理範例:

UseProxy=true;ProxyTransport=tcp;ProxyHost=socks.company.com;ProxyPort=1080;ProxySocksVersion=5

記錄設定

參數 類型 預設 Description
LogLevel 繩子 WARN 日誌層級:TRACE、DEBUG、INFO、WARN、ERROR

Example:

LogLevel=DEBUG

預設日誌位置:

${user.home}/.microsoft/livy-jdbc-driver/logs/driver.log

自訂日誌配置: 在你的 classpath 上使用自訂 log4j2.xml 檔案或 logback.xml 檔案。

高並行(HC)模式

高並行(HC)模式是為 JDBC 應用程式設計,這些應用程式能開啟多個連線至同一 Fabric 工作空間與湖庫。 Fabric 不再為每個連線配置獨立的經典 Spark 會話,而是以共享、伺服器管理的 Spark 容量來服務這些連線。

HC 模式能縮短連線啟動時間,避免重複 Spark 會話配置,並更有效率地利用可用容量以應付並行工作負載。 每個 JDBC 連線都會透過其在 HC 會話中指定的上下文持續執行語句。

HC 模式需選擇啟用,並可作為 JDBC 連線集區的補充。 HC 模式管理伺服器端的 Fabric Spark 會話,而 JDBC 連線池則管理用戶端應用程式中的連線。 如果你沒有啟用 HC 模式,驅動程式會建立一個經典的 Livy 會話。

在 HC 和經典模式之間選擇

以下情況下使用 HC 模式:

  • 應用程式服務處理同時執行 Spark SQL 的請求,這些請求會透過多個 JDBC 連線執行。
  • BI 或報表工具會對相同的工作區和 Lakehouse 開啟多個連線。
  • 並行的 ETL 任務或排程作業會造成 JDBC 連線暴增。
  • 連線使用相容的 Spark 配置,並可共享伺服器管理的 Spark 容量。
  • 減少連線啟動時間及重複配置 Spark 會話非常重要。

當以下情況使用經典模式:

  • 你的應用程式只使用一個或幾個長期存在的JDBC連線。
  • 每個連線都需要一個專用的 Spark 工作階段,以實現工作負載或資源的嚴格隔離。
  • 連接需要有實質差異的 Spark 配置或 Fabric 環境。
  • 目標工作空間和湖屋的 Fabric Livy 端點不支援 HC 模式。

啟用高並行模式

HC 模式需主動啟用。 如果你省略 hcEnabled 或設為 false,驅動程式會使用經典的 Livy 會話路徑。 設 hcEnabled 為 true 以取得 HC 會話:

hcEnabled=true;

HC 模式下的必要參數

啟用 HC 需要以下有效識別碼:

參數 註釋
WorkspaceId 必須是有效的 Fabric 工作區 GUID。
LakehouseId 必須是有效的 Fabric 湖倉 GUID。

包含 SessionTag 是因為它會影響伺服器端的熱場匹配。

會話匹配

HC 模式並不保證每個連線都會重複使用現有的會話。 對於每個會話擷取請求,驅動程式會發送:

  • SessionTag 值。
  • 有效的 Spark 配置,包括 conf.* 設定和 EnvironmentId。

Fabric Livy 服務會評估這些值,並決定是否將連線連接到預熱工作階段,或建立新的工作階段。 為了提升會話重用,應在旨在共享伺服器管理的 Spark 容量的連線間使用一致的數值。

關閉 HC 模式

你可以在不重新部署驅動程式的情況下,透過將 HC 連接字串 設為 false 或移除 HC 參數來停用新的 HC 會話擷取:

hcEnabled=false;

緊急回滾: 緊急回滾時,JVM 起始為 -Dcom.microsoft.fabric.jdbc.hc.disable=true。 這個 JVM 全域開關即使在連線字串包含 hcEnabled=true 時,也會強制使用經典模式。

了解 HC 執行時的行為

  • HC 的會話取得會被封鎖,直到會話準備好或 hcAcquireTimeoutSeconds 到期。
  • 在擷取過程中關閉連線會中斷擷取並觸發盡力而為的 HC 會話清理。
  • 在兩種模式下,身分驗證、查詢執行、取消、預備陳述式、結果處理以及 FabricEnvironmentID 都可透過標準 JDBC 介面運作。

使用範例

基本連接

import java.sql.*;

public class BasicConnectionExample {
    public static void main(String[] args) {
        String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
                     "FabricWorkspaceID=<workspace-id>;" +
                     "FabricLakehouseID=<lakehouse-id>;" +
                     "AuthFlow=2";
        
        try (Connection conn = DriverManager.getConnection(url)) {
            System.out.println("Connected successfully!");
            System.out.println("Database: " + conn.getMetaData().getDatabaseProductName());
            System.out.println("Driver: " + conn.getMetaData().getDriverName());
            System.out.println("Driver Version: " + conn.getMetaData().getDriverVersion());
        } catch (SQLException e) {
            System.err.println("Connection failed: " + e.getMessage());
            e.printStackTrace();
        }
    }
}

高並行模式連線

以下範例以 HC 模式開啟 JDBC 連線,執行查詢,並自動釋放連線:

String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=2;" +
             "hcEnabled=true;" +
             "sessionTag=dashboard_queries;" +
             "FabricEnvironmentID=<environment-id>";

try (Connection conn = DriverManager.getConnection(url);
     Statement stmt = conn.createStatement();
     ResultSet rs = stmt.executeQuery("SELECT current_timestamp()")) {
    while (rs.next()) {
        System.out.println(rs.getTimestamp(1));
    }
}

執行查詢

簡單查詢

public void executeSimpleQuery(Connection conn) throws SQLException {
    String sql = "SELECT current_timestamp() as now";
    
    try (Statement stmt = conn.createStatement();
         ResultSet rs = stmt.executeQuery(sql)) {
        
        if (rs.next()) {
            Timestamp now = rs.getTimestamp("now");
            System.out.println("Current timestamp: " + now);
        }
    }
}

使用篩選器的查詢

public void executeQueryWithFilter(Connection conn) throws SQLException {
    String sql = "SELECT * FROM sales WHERE amount > 1000 ORDER BY amount DESC";
    
    try (Statement stmt = conn.createStatement();
         ResultSet rs = stmt.executeQuery(sql)) {
        
        while (rs.next()) {
            int id = rs.getInt("id");
            double amount = rs.getDouble("amount");
            Date date = rs.getDate("sale_date");
            
            System.out.printf("ID: %d, Amount: %.2f, Date: %s%n", 
                            id, amount, date);
        }
    }
}

帶限制的查詢

public void executeQueryWithLimit(Connection conn) throws SQLException {
    String sql = "SELECT * FROM customers LIMIT 10";
    
    try (Statement stmt = conn.createStatement();
         ResultSet rs = stmt.executeQuery(sql)) {
        
        ResultSetMetaData metaData = rs.getMetaData();
        int columnCount = metaData.getColumnCount();
        
        // Print column names
        for (int i = 1; i <= columnCount; i++) {
            System.out.print(metaData.getColumnName(i) + "\t");
        }
        System.out.println();
        
        // Print rows
        while (rs.next()) {
            for (int i = 1; i <= columnCount; i++) {
                System.out.print(rs.getString(i) + "\t");
            }
            System.out.println();
        }
    }
}

結果集的處理

遍歷結果集

public void navigateResultSet(Connection conn) throws SQLException {
    String sql = "SELECT id, name, amount FROM orders";
    
    try (Statement stmt = conn.createStatement();
         ResultSet rs = stmt.executeQuery(sql)) {

        while (rs.next()) {
            System.out.printf(
                "Order %d: %s, %.2f%n",
                rs.getInt("id"),
                rs.getString("name"),
                rs.getDouble("amount"));
        }
    }
}

結果集是只能向前且唯讀的。 驅動程式不支援像 first()、 last()、 absolute()和 等方法。

處理大型結果集

驅動程式依回應大小分頁結果。 在 JDBC URL 中設定 LivyStatementPageSize,以設定要求的頁面大小(以位元組為單位)。 這個 Statement.setFetchSize() 方法不控制擷取或記憶體使用。

jdbc:fabricspark://api.fabric.microsoft.com;FabricWorkspaceID=<workspace-id>;FabricLakehouseID=<lakehouse-id>;AuthFlow=2;LivyStatementPageSize=1048576
public void processLargeResultSet(Connection conn) throws SQLException {
    String sql = "SELECT * FROM large_table";
    
    try (Statement stmt = conn.createStatement()) {
        try (ResultSet rs = stmt.executeQuery(sql)) {
            int rowCount = 0;
            while (rs.next()) {
                // Process row
                processRow(rs);
                rowCount++;
                
                if (rowCount % 10000 == 0) {
                    System.out.println("Processed " + rowCount + " rows");
                }
            }
            System.out.println("Total rows processed: " + rowCount);
        }
    }
}

private void processRow(ResultSet rs) throws SQLException {
    // Process individual row
}

使用預備陳述

public void usePreparedStatement(Connection conn) throws SQLException {
    String sql = "SELECT * FROM products WHERE category = ? AND price > ?";
    
    try (PreparedStatement pstmt = conn.prepareStatement(sql)) {
        // Set parameters
        pstmt.setString(1, "Electronics");
        pstmt.setDouble(2, 100.0);
        
        try (ResultSet rs = pstmt.executeQuery()) {
            while (rs.next()) {
                String name = rs.getString("name");
                double price = rs.getDouble("price");
                System.out.printf("Product: %s, Price: $%.2f%n", name, price);
            }
        }
    }
}

批次作業

這個驅動程式不支援 JDBC 交易。 批次執行是非交易性的,且部分成功。 後來的失敗不會推翻先前成功的說法。

public void executeBatchInsert(Connection conn) throws SQLException {
    String sql = "INSERT INTO logs (timestamp, level, message) VALUES (?, ?, ?)";
    int pendingStatements = 0;
    
    try (PreparedStatement pstmt = conn.prepareStatement(sql)) {
        // Add multiple statements to batch
        for (int i = 0; i < 1000; i++) {
            pstmt.setTimestamp(1, new Timestamp(System.currentTimeMillis()));
            pstmt.setString(2, "INFO");
            pstmt.setString(3, "Log message " + i);
            pstmt.addBatch();
            pendingStatements++;
            
            // Execute batch every 100 statements
            if ((i + 1) % 100 == 0) {
                pstmt.executeBatch();
                pstmt.clearBatch();
                pendingStatements = 0;
            }
        }
        
        if (pendingStatements > 0) {
            pstmt.executeBatch();
        }
        
        System.out.println("Batch insert completed successfully");
    }
}

選擇連線擁有者

選擇一個元件來負責連線生命週期。 不要將一個連線池置於另一個連線池之中。

使用量模型 入口 泳池擁有者 結清借來的債務 Connection 關機者
直接 DriverManager、 LivyDriver或 LivyDataSource None 實際關閉其獨立的 Livy 連線和工作階段。 應用程式會關閉所有連線。 直接 LivyDataSource 的生命週期沒有獨立的緊密生命週期。
外部 HikariCP HikariDataSource 並以 JDBC 網址為後盾或 LivyDataSource 僅限 HikariCP 將 Hikari 代理連線回傳給 HikariCP。 應用程式會關閉其保留的 HikariDataSource。
明確內建池 LivyBuiltInPooledDataSource 為每個已設定的資料來源產生器延遲建立一個 LivyConnectionPool 先解除邏輯租約,並且僅在安全清理完成後才重複使用實體工作階段。 應用程式關閉其保留的內建資料來源。
標準池化 SPI LivyConnectionPoolDataSource 為中階池管理器建立 PooledConnection 物件 中階的泳池經理,不是工廠 關閉邏輯控制代碼,並在完成清理後同步通知管理器。 管理器會呼叫 PooledConnection.close() 來銷毀每個實體連線。

ConnectionPoolEnabled 在驅動程式碼和隨附的屬性設定中皆預設為 false。 這並不是適用於目前入口點的使用者自行啟用開關。 如果應由驅動程式管理連線集區,請使用 LivyBuiltInPooledDataSource;或者設定外部連線集區,例如 HikariCP。 其他入口點則直接建立非擁有池的連線。

使用 HikariCP 的連線集區

HikariCP 可以透過 LivyDataSource 或透過 LivyDriver 中的 JDBC URL 建立連線。 建立一個 HikariCP 實例,重複使用該應用程式的整個壽命,並在任一配置中只使用一個連線池。

Maven 依賴

<dependency>
    <groupId>com.zaxxer</groupId>
    <artifactId>HikariCP</artifactId>
    <version>5.0.1</version>
</dependency>

使用 LivyDataSource

LivyDataSource 建立獨立的實體連線,且不擁有連線池或獨立的封閉生命週期。 直接供應給 HikariCP,讓 HikariCP 成為唯一的池擁有者:

import com.microsoft.spark.livy.jdbc.LivyDataSource;
import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;

import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;

public final class ConnectionPoolExample implements AutoCloseable {
    private final HikariDataSource hikariDataSource;

    public ConnectionPoolExample(String url) {
        LivyDataSource configuredDataSource = new LivyDataSource();
        configuredDataSource.setUrl(url);

        HikariConfig config = new HikariConfig();
        config.setDataSource(configuredDataSource);
        config.setMaximumPoolSize(2);
        config.setMinimumIdle(0);
        config.setConnectionTimeout(900000);     // Wait up to 15 minutes for a connection
        config.setInitializationFailTimeout(-1); // Connect on demand
        config.setIdleTimeout(600000);           // 10 minutes
        config.setMaxLifetime(1800000);          // 30 minutes
        config.setPoolName("FabricSparkPool");

        hikariDataSource = new HikariDataSource(config);
    }

    public Connection getConnection() throws SQLException {
        return hikariDataSource.getConnection();
    }

    @Override
    public void close() {
        hikariDataSource.close();
    }

    public static void main(String[] args) throws SQLException {
        String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
                     "FabricWorkspaceID=<workspace-id>;" +
                     "FabricLakehouseID=<lakehouse-id>;" +
                     "AuthFlow=AZURE_CLI;" +
                     "hcEnabled=true;" +
                     "sessionTag=pooled_app;" +
                     "hcAcquireTimeoutSeconds=600";

        try (ConnectionPoolExample pool = new ConnectionPoolExample(url);
             Connection conn = pool.getConnection();
             Statement stmt = conn.createStatement();
             ResultSet rs = stmt.executeQuery("SELECT 'Pooled connection!' as message")) {

            if (rs.next()) {
                System.out.println(rs.getString("message"));
            }
        }
    }
}

範例中啟用了 JDBC URL 中的 HC 模式。 HikariCP 管理用戶端連線,而 HC 模式則管理 Fabric 中的共享 Spark 會話。 如果應用程式需要經典會話,請移除 hcEnabled 和 sessionTag 。

直接使用LivyDriver

設定 JDBC URL 與驅動程式類別,讓 HikariCP 透過 LivyDriver 建立實體連線。 當可使用 JDBC 服務提供者自動探索機制時,可選擇是否明確指定驅動程式類別。 此配置不會建立 LivyDataSource 的實例,因此其內建集區不會涉及:

import com.zaxxer.hikari.HikariConfig;
import com.zaxxer.hikari.HikariDataSource;

import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;

public final class DriverConnectionPoolExample {
    public static HikariDataSource createPool(String url) {
        HikariConfig config = new HikariConfig();
        config.setDriverClassName("com.microsoft.spark.livy.jdbc.LivyDriver");
        config.setJdbcUrl(url);

        // Keep pool sizes small: each classic-mode physical connection owns a Livy session.
        config.setMaximumPoolSize(2);
        config.setMinimumIdle(0);
        config.setConnectionTimeout(900000);     // Wait up to 15 minutes to borrow a connection
        config.setInitializationFailTimeout(-1); // Skip startup validation; connect on demand
        config.setIdleTimeout(600000);           // 10 minutes
        config.setMaxLifetime(1800000);          // 30 minutes
        config.setPoolName("FabricSparkDriverPool");

        return new HikariDataSource(config);
    }

    public static void main(String[] args) throws SQLException {
        String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
                     "FabricWorkspaceID=<workspace-id>;" +
                     "FabricLakehouseID=<lakehouse-id>;" +
                     "AuthFlow=AZURE_CLI;" +            // Uses the DefaultAzureCredential chain
                     "hcEnabled=true;" +
                     "sessionTag=pooled_app;" +
                     "hcAcquireTimeoutSeconds=600";  // Poll up to 10 minutes for HC readiness

        // Reuse this pool for the application lifetime; main closes it at process exit.
        try (HikariDataSource pool = createPool(url);
             Connection conn = pool.getConnection();
             Statement stmt = conn.createStatement();
             ResultSet rs = stmt.executeQuery("SELECT 'Pooled LivyDriver connection!' as message")) {

            if (rs.next()) {
                System.out.println(rs.getString("message"));
            }
        }
    }
}

設定池時請參考以下指引:

設定或行為 Guidance
游泳池所有權 LivyDataSource 僅限直接撥號。 當你將其提供給 HikariCP 時,該連線池就會由 HikariCP 接管。 當驅動程式應擁有該集區時,請改用 LivyBuiltInPooledDataSource。
集區大小 先從一個小池開始,然後在測量工作負載並發性和容量使用後再擴大。 在經典模式下,每個實體連線都擁有一個 Livy 會話。 在 HC 模式下,連線可以共享伺服器管理的 Spark 容量。
connectionTimeout 控制來電者等待借用連線的時間,但不會取消正在進行的連線嘗試。 對於 HC URL,可使用 hcAcquireTimeoutSeconds;在經典模式下,可使用 LivySessionTimeoutSeconds、驗證、HTTP 重試和驗證檢查。 HikariCP 以毫秒級表達此設定。
initializationFailTimeout 負值則跳過啟動連線與驗證嘗試。 連接失敗會在應用程式首次請求連線時浮現。 若應用程式必須在啟動時驗證連接性,則使用正值。
連線驗證 HikariCP 使用 Connection.isValid()。 其有效持續時間由驅動程式的 HTTP 逾時及重試/退回設定控制,而非 HikariCP 的 validationTimeout。 不要設定連線測試查詢,因為它會用 Spark 語句取代較輕的驗證檢查。

完整的 HikariCP 配置選項列表,請參見:

資料型態映射

驅動程式將 Spark SQL 資料型態映射到 JDBC SQL 型態與 Java 型別:

Spark SQL 類型 JDBC SQL 類型 Java 類型 註釋
BOOLEAN BOOLEAN Boolean
BYTE TINYINT Byte
SHORT SMALLINT Short
INT INTEGER Integer
LONG BIGINT Long
FLOAT FLOAT Float
DOUBLE DOUBLE Double
DECIMAL DECIMAL BigDecimal 精確度與規模得以保存
STRING VARCHAR String
VARCHAR(n) VARCHAR String
CHAR(n) CHAR String
BINARY BINARY byte[]
DATE DATE java.sql.Date
TIMESTAMP TIMESTAMP java.sql.Timestamp
ARRAY ARRAY String 以 JSON 文字擷取;中繼資料類型不保證對 java.sql.Array 存取器的完整支援
MAP JAVA_OBJECT String 以 JSON 文字擷取
STRUCT STRUCT String 以 JSON 文字檢索;元資料類型不保證完整的 java.sql.Struct 存取器支援