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.
- Verwendung von DBeaver(JDBC clients)
- Verwendung des Codes Java(JDBC Client)
- Verwendung von Python(PyHive JDBC Client)
Vorbereitende Schritte
- Installieren Sie watsonx.data.
- Bereitstellung der nativen Spark-Engine in watsonx.data.
- Laden Sie den JDBC Client herunter:
queryserver-jdbc-SNAPSHOT-standalone.jarüber den Download-Link. - Führen Sie den Spark-Query-Server in der Spark-Engine aus. Um einen neuen Abfrageserver zu erstellen, siehe Erstellen eines Spark-Abfrageservers.
- 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.
-
Öffnen Sie DBeaver und klicken Sie in der Menüleiste auf Datenbank > Treibermanager.
-
Suche nach Hive. Sie finden Apache Hive 4+ Treiber unter Hadoop kategorie.
-
Klicken Sie auf Kopieren.
-
Ändern Sie den Namen in Spark watsonx.data.
-
Ändern Sie die folgenden Einstellungen:
-
Auf der Registerkarte Einstellungen,
-
Auf der Registerkarte Bibliotheken- Fügen Sie die Spark JDBC Query Server JAR-Datei hinzu.
-
Wählen Sie Datenbank Navigator, klicken Sie auf Neue Verbindung und führen Sie die folgenden Schritte aus:
- Wählen Sie den neu erstellten Treiber aus.
- Klicken Sie auf Verbinden durch und wählen Sie URL.
- Geben Sie die JDBC URL in folgendem Format an:
jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>. - Wählen Sie Authentifizierung, geben Sie Username als Ihren Benutzernamen und Ihren IAM-API-Schlüssel als Passwort an.
- 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:
-
Vergewissern Sie sich, dass Sie Python Version 3.12 oder niedriger haben.
-
Installieren Sie pyHive mit
pip install thrift "PyHive[hive_pure_sasl]==0.7.0". -
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() -
Starten Sie mit
python connect.py.