Verbindung zum Spark-Abfrageserver mit dem Spark-Treiber JDBC

Gilt für: Funkenmotor

Sie können sich auf folgende Weise mit dem Spark-Abfrageserver verbinden und Abfragen zur Analyse Ihrer Daten ausführen.

Vorbereitende Schritte

  1. Installieren Sie watsonx.data.
  2. Bereitstellung der nativen Spark-Engine in watsonx.data.
  3. Laden Sie den JDBC Client herunter: queryserver-jdbc-SNAPSHOT-standalone.jar über den Download-Link.
  4. Führen Sie den Spark-Query-Server in der Spark-Engine aus. Um einen neuen Abfrageserver zu erstellen, siehe Erstellen eines Spark-Abfrageservers.
  5. Verbindungseigenschaften - Klicken Sie auf das Drei-Punkte-Menü für Query Server, klicken Sie auf Verbindungsdetails und kopieren Sie die folgenden Verbindungsdetails:
    • Host
    • URI
    • Instanz
    • Benutzername
    • Ihr IBM IAM API-Schlüssel.

Verbindung zum Spark-Abfrageserver mit Hilfe von DBeaver ( JDBC client)

Um eine Verbindung zum Spark-Abfrageserver mit einem JDBC-Client, wie z. B. DBeaver, herzustellen, richten Sie den watsonx.data-Treiber in DBeaver ein.

  1. Öffnen Sie DBeaver und klicken Sie in der Menüleiste auf Datenbank > Treibermanager.

  2. Suche nach Hive. Sie finden Apache Hive 4+ Treiber unter Hadoop kategorie.

  3. Klicken Sie auf Kopieren.

  4. Ändern Sie den Namen in Spark watsonx.data.

  5. Ändern Sie die folgenden Einstellungen:

  6. Auf der Registerkarte Einstellungen,

  7. Auf der Registerkarte Bibliotheken- Fügen Sie die Spark JDBC Query Server JAR-Datei hinzu.

  8. Wählen Sie Datenbank Navigator, klicken Sie auf Neue Verbindung und führen Sie die folgenden Schritte aus:

    1. Wählen Sie den neu erstellten Treiber aus.
    2. Klicken Sie auf Verbinden durch und wählen Sie URL.
    3. Geben Sie die JDBC URL in folgendem Format an: jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>.
    4. Wählen Sie Authentifizierung, geben Sie Username als Ihren Benutzernamen und Ihren IAM-API-Schlüssel als Passwort an.
    5. Speichern und verbinden Sie sich mit der Verbindung durch Doppelklick.

Verbindung zum Spark-Abfrageserver mit Hilfe von Java ( JDBC Client) Code

Stellen Sie sicher, dass Ihr Java CLASSPATH den heruntergeladenen JDBC Treiber enthält. Zum Beispiel:

java -cp queryserver-jdbc-SNAPSHOT-standalone.jar App.java

Sie können die Parameterwerte angeben und den folgenden Java Code verwenden, um eine Verbindung mit dem Spark-Abfrageserver herzustellen. Wenn Sie die API v2 verwenden, setzen Sie den Parameter <api_version> auf v2; für die API v3 setzen Sie ihn auf v3.

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.ResultSetMetaData;
import java.sql.Statement;

public class App {
    public static void main(String[] args) throws Exception {
        // Set the below configurations from Connection Details of QueryS Server
        // Exclude having https/http/www, just domain
        String host = "example.com";
        String Instance = "CRN/OR/INSTANCE-ID";
        String uri = "/lakehouse/api/<api_version>/spark_engines/.../query_servers/.../connect/cliservice";
        String user = "EMAIL-ID/OR/USER-ID";
        String apikey = "API-KEY";

        String jdbcUrl = String.format("jdbc:hive2://%s/default;instance=%s;httpPath=%s;", host, Instance, uri);

        // Required if your domain requires SSL certificates
        // This is not required for SaaS, hence comment the below line for SaaS
        // Else, we need provide trust-store path which has the SSL certificates for the host
        jdbcUrl += "sslTrustStore=tech_trust.jks;trustStorePassword=Test@123";

        try {
            // Load the Hive JDBC driver
            Class.forName("com.ibm.wxd.spark.jdbc.QueryServerDriver");

            // Connect to Hive
            Connection con = DriverManager.getConnection(jdbcUrl, user, apikey);
            Statement stmt = con.createStatement();

            System.out.println("Connected to watsonx.data Spark Query Server");

            // Sample query
            String sql = "show databases";

            ResultSet rs = stmt.executeQuery(sql);
            ResultSetMetaData rsmd = rs.getMetaData();
            int columnCount = rsmd.getColumnCount();

            // The column count starts from 1
            for (int i = 1; i <= columnCount; i++ ) {
                System.out.println(rsmd.getColumnName(i));
            }

            // Print result
            while (rs.next()) {
                System.out.println(rs.getString(1)); // Or loop through columns
            }

            // Clean up
            rs.close();
            stmt.close();
            con.close();

        } catch (Exception e) {
            e.printStackTrace();
        }
    }

Verbindung zum Spark-Abfrageserver mit Hilfe von Python ( PyHive JDBC Client)

Um eine Verbindung zum Spark-Abfrageserver mit einem Programm Python herzustellen, gehen Sie wie folgt vor:

  1. Vergewissern Sie sich, dass Sie Python Version 3.12 oder niedriger haben.

  2. Installieren Sie pyHive mit pip install thrift "PyHive[hive_pure_sasl]==0.7.0".

  3. Speichern Sie das Ergebnis in einer Datei wie connect.py.

    
    import ssl
    import thrift
    import base64
    from pyhive import hive
    
    import requests
    import thrift.transport
    import thrift.transport.THttpClient
    
    import logging
    import contextlib
    from http.client import HTTPConnection
    
    
    # Change the following inputs. When using the v2 API, set the <api_version> parameter to `v2`; for the v3 API, set it to `v3`.
    class Credentials:
        host = "https://example.ibm.com"
        uri = "/lakehouse/api/<api_version>/spark_engines/.../query_servers/.../connect/cliservice"
        instance_id = "CRN/OR/INSTANCE-ID"
        username = "EMAIL-ID/OR/USER-ID"
        apikey = "API-KEY"
    
    
    creds = Credentials()
    
    
    def disable_ssl(ctx):
        ctx.check_hostname = False
        ctx.verify_mode = ssl.CERT_NONE
    
        ssl.SSLContext.verify_mode = property(lambda self: ssl.CERT_NONE, lambda self, newval: None)
    
    
    def get_access_token(apikey):
        try:
            headers = {
                'Content-Type': 'application/x-www-form-urlencoded',
                'Accept': 'application/json',
            }
    
            data = {
                'grant_type': 'urn:ibm:params:oauth:grant-type:apikey',
                'apikey': apikey,
            }
    
            response = requests.post('https://iam.cloud.ibm.com/identity/token', headers=headers, data=data)
            return response.json()['access_token']
        except Exception as inst:
            print('Error in getting access token')
            print(inst)
            exit
    
    ctx = ssl.create_default_context()
    
    ## If you require to disable SSL, uncomment the below line
    # disable_ssl(ctx)
    
    transport = thrift.transport.THttpClient.THttpClient(
        uri_or_host="{host}:{port}{uri}".format(
            host=creds.host, uri= creds.uri, port=443,
        ),
        ssl_context=ctx,
    )
    
    headers = {
        "AuthInstanceId": creds.instance_id
    }
    
    if creds.instance_id.isdigit():
        # Software installation
        headers["Authorization"] =  "ZenApiKey " + base64.b64encode(f"{creds.username}:{creds.apikey}".encode('utf-8')).decode('utf-8')
    else:
        # Cloud installation
        headers["Authorization"] = "Bearer {}".format(get_access_token(creds.apikey))
    
    transport.setCustomHeaders(headers)
    
    cursor = hive.connect(thrift_transport=transport).cursor()
    print("Connected to Spark Query Server")
    
    cursor.execute('show databases')
    print(cursor.fetchall())
    
    cursor.close()
    
    
  4. Starten Sie mit python connect.py.