sterownik Microsoft JDBC na potrzeby inżynierii danych Microsoft Fabric

JDBC (łączność z bazą danych Java) to powszechnie przyjęty standard, który umożliwia aplikacjom klienckim łączenie się z danymi z baz danych i platform danych big data oraz pracę z nimi.

Microsoft JDBC Driver for Fabric Data Engineering pozwala łączyć, zapytywać i zarządzać obciążeniami Spark w Fabric z niezawodnością i prostotą standardu JDBC. Oparty na API Livy firmy Fabric, sterownik zapewnia bezpieczną i elastyczną łączność Spark SQL z aplikacjami Java i narzędziami BI. Ta integracja umożliwia bezpośrednie przesyłanie i wykonywanie kodu Spark bez konieczności tworzenia oddzielnych elementów notesu lub elementów definicji zadania Spark. Sterownik jest zgodny z popularnymi klientami JDBC, takimi jak DbVisualizer i DBeaver, a także narzędziami analizy biznesowej obsługującymi łączność JDBC, w tym Tableau.

Najważniejsze funkcje    

  • Zgodność JDBC 4.2: pełna implementacja specyfikacji JDBC 4.2
  • Uwierzytelnianie Microsoft Entra ID: wiele przepływów uwierzytelniania, w tym interakcyjne, poświadczenia klienta i uwierzytelnianie oparte na certyfikatach
  • Integracja z HikariCP: Użyj HikariCP do zarządzania i ponownego wykorzystania połączeń JDBC dla aplikacji produkcyjnych
  • Obsługa zapytań natywnych spark SQL: bezpośrednie wykonywanie instrukcji Spark SQL bez tłumaczenia
  • Kompleksowa obsługa typów danych: obsługa wszystkich typów danych Spark SQL, w tym typów złożonych (ARRAY, MAP, STRUCT)
  • Asynchroniczne pobieranie zestawu wyników: ładowanie danych w tle w celu zwiększenia wydajności
  • Wzorzec obwodu przerywającego: Ochrona przed awariami kaskadowymi z automatycznym ponawianiem
  • Automatyczne ponowne nawiązywanie połączenia: Przezroczyste odzyskiwanie sesji w przypadku niepowodzeń połączenia
  • Zaawansowana logika ponawiania prób: ponawianie prób z wykładniczym opóźnieniem i odzyskiwanie sesji w celu zwiększenia odporności.
  • Obsługa serwera proxy: konfiguracja serwera proxy HTTP i SOCKS dla środowisk przedsiębiorstwa

Wymagania wstępne

Przed użyciem sterownika Microsoft JDBC dla Fabric Data Engineering, upewnij się, że masz:

  • Zestaw Java Development Kit (JDK): wersja 11 lub nowsza (zalecane środowisko Java 21)
  • Fabric Access: Dostęp do przestrzeni roboczej Fabric
  • Microsoft Entra ID poświadczenia: Odpowiednie poświadczenia do uwierzytelniania
  • Identyfikatory obszaru roboczego i lakehouse: identyfikatory GUID dla obszaru roboczego Fabric i lakehouse

Pobieranie i instalacja

Microsoft JDBC Driver for Fabric Data Engineering w wersji 1.0.0 obsługuje Java 11, 17 i 21. Stale ulepszamy obsługę łączności w języku Java i zalecamy pracę z najnowszą wersją sterownika JDBC firmy Microsoft.

  1. Pobierz plik zip lub tar z powyższych linków.
  2. Wyodrębnij pobrany plik, aby uzyskać dostęp do plików JAR sterownika.
  3. Wybierz plik JAR zgodny z wersją środowiska JRE:
    • W przypadku języka Java 11: ms-sparksql-jdbc-1.0.1.jre11.jar
    • Dla środowiska Java 17: ms-sparksql-jdbc-1.0.1.jre17.jar
    • Dla środowiska Java 21: ms-sparksql-jdbc-1.0.1.jre21.jar
  4. Dodaj wybrany plik JAR do ścieżki klasy aplikacji.
  5. W przypadku klientów JDBC skonfiguruj klasę sterowników JDBC: com.microsoft.spark.livy.jdbc.LivyDriver

Przykład szybkiego startu

Ten przykład pokazuje, jak połączyć się z Fabric i wykonać zapytanie za pomocą sterownika Microsoft JDBC dla Fabric Data Engineering. Przed uruchomieniem tego kodu upewnij się, że zostały spełnione wymagania wstępne i zainstalowano sterownik.

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;" +  // Azure CLI based authentication
                     "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();
        }
    }
}

Format ciągu połączenia

Podstawowy ciąg połączenia

Microsoft JDBC Driver for Fabric Data Engineering używa następującego formatu parametrów połączenia:

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

Składniki ciągu połączenia

Składnik Description Example
Protocol Identyfikator protokołu URL JDBC jdbc:fabricspark://
Nazwa hosta Nazwa hosta punktu końcowego usługi Fabric api.fabric.microsoft.com
Port Opcjonalny numer portu (domyślnie: 443) :443
Parametry Rozdzielane średnikami pary klucz=wartość FabricWorkspaceID=<guid>

Przykładowe ciągi połączenia

Połączenie podstawowe (uwierzytelnianie oparte na przeglądarce interakcyjnej)

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

Za pomocą konfiguracji zasobów platformy 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

Za pomocą właściwości sesji Spark

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

Authentication

Sterownik Microsoft JDBC dla Fabric Data Engineering obsługuje wiele metod uwierzytelniania poprzez Microsoft Entra ID (dawniej Azure Active Directory). Uwierzytelnianie jest konfigurowane przy użyciu parametru AuthFlow w parametrach połączenia.

Przepływy uwierzytelniania

AuthFlow Metoda uwierzytelniania Przypadek użycia
1 Przeglądarka interaktywna Uwierzytelnianie interakcyjne użytkownika przy użyciu protokołu OAuth 2.0
2 Azure CLI Programowanie przy użyciu interfejsu wiersza polecenia platformy Azure
3 Poświadczenia Sekretne Klienta (Podstawowy Element Usługi) Automatyczne uwierzytelnianie między usługami
4 Certyfikat klienta - poświadczenie Uwierzytelnianie jednostki usługi opartej na certyfikatach
5 Token dostępu Wstępnie uzyskany token dostępu nosiciela

Uwierzytelnianie za pomocą przeglądarki interakcyjnej

Najlepsze do: Rozwoju i interaktywnych aplikacji

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

Connection conn = DriverManager.getConnection(url);

Parametry:

  • AuthFlow=1: Określa uwierzytelnianie za pomocą przeglądarki interakcyjnej
  • AuthTenantID (opcjonalnie): identyfikator dzierżawy Microsoft Entra
  • AuthClientID (opcjonalnie): identyfikator aplikacji (klienta)

Zachowanie:

  • Otwiera okno przeglądarki na potrzeby uwierzytelniania użytkownika
  • Poświadczenia są buforowane dla kolejnych połączeń do momentu wygaśnięcia
  • Odpowiednie dla aplikacji z jednym użytkownikiem

Poświadczenia klienta lub uwierzytelnianie jednostki usługi

Najlepsze dla: Zautomatyzowane usługi i zadania w tle

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);

Wymagane parametry:

  • AuthFlow=3: Określa uwierzytelnianie poświadczeń klienta
  • AuthClientID: Identyfikator aplikacji (klienta) z Microsoft Entra ID
  • AuthClientSecret: Tajny klucz klienta od Microsoft Entra ID
  • AuthTenantID: identyfikator dzierżawy Microsoft Entra

Najlepsze rozwiązania:

  • Bezpieczne przechowywanie wpisów tajnych (azure Key Vault, zmienne środowiskowe)
  • Używanie tożsamości zarządzanych, jeśli jest to możliwe
  • Regularne zmienianie sekretów

Uwierzytelnianie certyfikatem

Najlepsze rozwiązanie: Aplikacje dla przedsiębiorstw wymagające uwierzytelniania opartego na certyfikatach

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

Connection conn = DriverManager.getConnection(url);

Wymagane parametry:

  • AuthFlow=4: Określa uwierzytelnianie oparte na certyfikatach
  • AuthClientID: Identyfikator aplikacji (klienta)
  • AuthCertificatePath: Ścieżka do pliku certyfikatu PFX/PKCS12
  • AuthCertificatePassword: Hasło certyfikatu
  • AuthTenantID: identyfikator dzierżawy Microsoft Entra

Uwierzytelnianie tokenu dostępu

Najlepszy wybór dla: niestandardowe scenariusze uwierzytelniania

// 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);

Buforowanie uwierzytelniania

Sterownik automatycznie buforuje tokeny uwierzytelniania w celu zwiększenia wydajności:

// Enable/disable caching (enabled by default)
String url = "jdbc:fabricspark://api.fabric.microsoft.com;" +
             "FabricWorkspaceID=<workspace-id>;" +
             "FabricLakehouseID=<lakehouse-id>;" +
             "AuthFlow=2;" +
             "AuthEnableCaching=true;" +  // Enable token caching
             "AuthCacheTTLMS=3600000";    // Cache TTL: 1 hour

Connection conn = DriverManager.getConnection(url);

Parametry konfiguracji

Wymagane parametry

Te parametry muszą być obecne w każdym ciągu połączenia:

Parameter Typ Description Example
FabricWorkspaceID UUID Identyfikator obszaru roboczego Fabric <workspace-id>
FabricLakehouseID UUID Identyfikator jeziora Fabric <lakehouse-id>
AuthFlow Integer Typ przepływu uwierzytelniania (1–5) 2

Parametry opcjonalne

Konfiguracja wersji interfejsu API

Parameter Typ Default Description
FabricVersion Sznurek v1 Wersja Fabric API
LivyApiVersion Sznurek 2023-12-01 Wersja interfejsu API usługi Livy

Konfiguracja środowiska

Parameter Typ Default Description
FabricEnvironmentID UUID Żaden Identyfikator środowiska Fabric do odwoływania się do składnika środowiska dla sesji Spark

Konfiguracja platformy Spark

Konfiguracja zasobu sesji

Skonfiguruj zasoby sesji platformy Spark w celu uzyskania optymalnej wydajności:

Parameter Typ Default Description Example
DriverCores Integer Ustawienie domyślne platformy Spark Liczba rdzeni CPU dla sterownika 4
DriverMemory Sznurek Ustawienie domyślne platformy Spark Alokacja pamięci dla sterownika 4g
ExecutorCores Integer Ustawienie domyślne platformy Spark Liczba rdzeni procesora CPU na egzekutora 4
ExecutorMemory Sznurek Ustawienie domyślne platformy Spark Alokacja pamięci na wykonawcę 8g
NumExecutors Integer Ustawienie domyślne platformy Spark Liczba wykonawców 2

Example:

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

Niestandardowe właściwości sesji platformy Spark

Dowolny parametr z prefiksem spark. jest automatycznie stosowany do sesji platformy Spark:

Przykładowe konfiguracje platformy 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

Natywny silnik wykonywania (NEE):

spark.nee.enabled=true

Kompletny przykład:

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

Ustawienia połączenia klienta HTTP

Konfiguruj połączenia transportowe HTTP sterownika dla optymalnej wydajności sieci. Te ustawienia nie konfigurują ani nie zarządzają pulą połączeń JDBC:

Parameter Typ Default Description
HttpMaxTotalConnections Integer 100 Maksymalna łączna liczba połączeń HTTP
HttpMaxConnectionsPerRoute Integer 50 Maksymalna liczba połączeń na trasę
HttpConnectionTimeoutInSeconds Integer 30 Przekroczenie limitu czasu połączenia
HttpSocketTimeoutInSeconds Integer 60 Limit czasu odczytu gniazda
HttpReadTimeoutInSeconds Integer 60 Limit czasu oczekiwania na odczyt HTTP
HttpConnectionRequestTimeoutSeconds Integer 30 Przekroczenie czasu żądania połączenia z puli połączeń
HttpEnableKeepAlive logiczny true Włączanie utrzymywania aktywności protokołu HTTP
HttpKeepAliveTimeoutSeconds Integer 60 Limit czasu podtrzymania połączenia
HttpFollowRedirects logiczny true Śledź przekierowania HTTP
HttpUseAsyncIO logiczny false Używanie asynchronicznego I/O HTTP

Example:

HttpMaxTotalConnections=200;HttpMaxConnectionsPerRoute=100;HttpConnectionTimeoutInSeconds=60

Konfiguracja serwera proxy

Skonfiguruj ustawienia serwera proxy HTTP i SOCKS dla środowisk przedsiębiorstwa:

Parameter Typ Default Description
UseProxy logiczny false Włączanie serwera proxy
ProxyTransport Sznurek http Typ transportu serwera proxy (http/tcp)
ProxyHost Sznurek Żaden Nazwa hosta serwera proxy
ProxyPort Integer Żaden Port serwera proxy
ProxyAuthEnabled logiczny false Włączanie uwierzytelniania serwera proxy
ProxyUsername Sznurek Żaden Nazwa użytkownika uwierzytelniania serwera proxy
ProxyPassword Sznurek Żaden Hasło uwierzytelniania serwera proxy
ProxyAuthScheme Sznurek basic Schemat uwierzytelniania (basic/digest/ntlm)
ProxySocksVersion Integer 5 Wersja protokołu SOCKS (4/5)

Przykład serwera proxy HTTP:

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

Przykład serwera proxy SOCKS:

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

Konfiguracja rejestrowania

Parameter Typ Default Description
LogLevel Sznurek INFO Poziom rejestrowania: TRACE (śledzenie), DEBUG (debugowanie), INFO (informacja), WARN (ostrzeżenie), ERROR (błąd)

Example:

LogLevel=DEBUG

Domyślna lokalizacja dziennika:

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

Niestandardowa konfiguracja dziennika: Użyj niestandardowego pliku lub log4j2.xmllogback.xml na ścieżce klas.


Przykłady użycia

Połączenie podstawowe

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();
        }
    }
}

Wykonywanie zapytań

Proste zapytanie

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);
        }
    }
}

Zapytanie z filtrem

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);
        }
    }
}

Wykonywanie zapytań z limitem

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();
        }
    }
}

Praca z zestawami wyników

public void navigateResultSet(Connection conn) throws SQLException {
    String sql = "SELECT id, name, amount FROM orders";
    
    try (Statement stmt = conn.createStatement(
            ResultSet.TYPE_SCROLL_INSENSITIVE,
            ResultSet.CONCUR_READ_ONLY);
         ResultSet rs = stmt.executeQuery(sql)) {
        
        // Move to first row
        if (rs.first()) {
            System.out.println("First row: " + rs.getString("name"));
        }
        
        // Move to last row
        if (rs.last()) {
            System.out.println("Last row: " + rs.getString("name"));
            System.out.println("Total rows: " + rs.getRow());
        }
        
        // Move to specific row
        if (rs.absolute(5)) {
            System.out.println("Row 5: " + rs.getString("name"));
        }
    }
}

Przetwarzanie dużych zestawów wyników

public void processLargeResultSet(Connection conn) throws SQLException {
    String sql = "SELECT * FROM large_table";
    
    try (Statement stmt = conn.createStatement()) {
        // Set fetch size for efficient memory usage
        stmt.setFetchSize(1000);
        
        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
}

Korzystanie z przygotowanych zapytań

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);
            }
        }
    }
}

Operacje wsadowe

public void executeBatchInsert(Connection conn) throws SQLException {
    String sql = "INSERT INTO logs (timestamp, level, message) VALUES (?, ?, ?)";
    
    try (PreparedStatement pstmt = conn.prepareStatement(sql)) {
        conn.setAutoCommit(false);  // Disable auto-commit for batch
        
        // 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();
            
            // Execute batch every 100 statements
            if (i % 100 == 0) {
                pstmt.executeBatch();
                pstmt.clearBatch();
            }
        }
        
        // Execute remaining statements
        pstmt.executeBatch();
        conn.commit();
        
        System.out.println("Batch insert completed successfully");
    } catch (SQLException e) {
        conn.rollback();
        throw e;
    } finally {
        conn.setAutoCommit(true);
    }
}

Pulowanie połączeń z HikariCP

Skonfiguruj HikariCP z adresem URL JDBC, aby tworzył połączenia przez LivyDriver. Stwórz jedną instancję HikariCP i wykorzystaj ją ponownie przez cały okres życia aplikacji.

Zależność Maven

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

Skonfiguruj HikariCP z URL-em JDBC

Skonfiguruj adres URL JDBC i klasę sterownika tak, aby HikariCP tworzył fizyczne połączenia przez LivyDriver. Ustawienie klasy sterownika jawnie jest opcjonalne, gdy dostępne jest automatyczne wykrywanie przez dostawcę usług JDBC:

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
                     "LivySessionTimeoutSeconds=600";  // Poll up to 10 minutes for session 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"));
            }
        }
    }
}

Skorzystaj z następujących wskazówek przy konfigurowaniu puli:

Ustawienie lub zachowanie Guidance
connectionTimeout Określa, jak długo wywołujący czeka na uzyskanie połączenia, ale nie anuluje już trwającej próby nawiązania połączenia. Ustaw ją powyżej, LivySessionTimeoutSeconds aby zapewnić czas na uwierzytelnianie, tworzenie sesji, próby HTTP i początkową weryfikację.
Walidacja połączeń HikariCP używa Connection.isValid(). O jego efektywnym czasie trwania decydują limit czasu HTTP sterownika oraz ustawienia ponawiania prób i opóźnienia po nieudanej próbie (backoff), a nie ustawienie validationTimeout w HikariCP. Nie konfiguruj zapytania testowego połączenia, ponieważ zastępuje ono mniej obciążające sprawdzenie poprawności instrukcją języka Spark.

Pełną listę opcji konfiguracji HikariCP można znaleźć tutaj:

Mapowanie typów danych

Sterownik mapuje typy danych Spark SQL na typy SQL JDBC i typy języka Java:

Typ Spark SQL Typ SQL JDBC Typ języka Java Notatki
BOOLEAN BOOLEAN Boolean
BYTE TINYINT Byte
SHORT SMALLINT Short
INT INTEGER Integer
LONG BIGINT Long
FLOAT FLOAT Float
DOUBLE DOUBLE Double
DECIMAL DECIMAL BigDecimal Precyzja i skalowanie zachowane
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 VARCHAR String Serializowany jako kod JSON
MAP VARCHAR String Serializowany jako kod JSON
STRUCT VARCHAR String Serializowany jako kod JSON