Connexion au serveur de requêtes Spark en utilisant le pilote Spark JDBC

S'applique à: Moteur à étincelles

Vous pouvez vous connecter au serveur de requêtes Spark de la manière suivante et exécuter des requêtes pour analyser vos données.

Avant de commencer

  1. Installer watsonx.data.
  2. Provisionner le moteur Spark natif dans watsonx.data.
  3. Téléchargez le client JDBC: queryserver-jdbc-SNAPSHOT-standalone.jar à partir du lien de téléchargement.
  4. Exécutez le serveur de requêtes Spark dans le moteur Spark. Pour créer un nouveau serveur de requêtes, voir Créer un serveur de requêtes Spark.
  5. Propriétés de la connexion - Cliquez sur le menu à trois points pour Query Server, cliquez sur Détails de la connexion et copiez les détails de la connexion suivants :
    • Hôte
    • URI
    • Instance
    • Nom d’utilisateur
    • Votre clé d'API IBM IAM.

Connexion au serveur de requêtes Spark en utilisant DBeaver ( JDBC client)

Pour se connecter au serveur de requêtes Spark à l'aide d'un client JDBC, tel que DBeaver, configurez le pilote watsonx.data dans DBeaver.

  1. Ouvrez DBeaver et dans la barre de menu cliquez sur Database > Driver Manager.

  2. Recherche de Hive. Vous pouvez trouver le pilote Apache Hive 4+ sous Hadoop catégorie.

  3. Cliquez sur Copier.

  4. Changez le nom en Spark watsonx.data.

  5. Modifier les paramètres suivants :

  6. Dans l'onglet Paramètres,

  7. Dans l'onglet Libraries- Ajouter le fichier JAR Spark JDBC query server.

  8. Sélectionnez Database Navigator, cliquez sur New Connection et effectuez les étapes suivantes :

    1. Sélectionnez le pilote nouvellement créé.
    2. Cliquez sur Connecter par et sélectionnez URL.
    3. Fournissez le JDBC URL en utilisant le format suivant : jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>.
    4. Sélectionnez Authentification, indiquez Username comme nom d'utilisateur et votre clé IAM API comme mot de passe.
    5. Enregistrez et connectez-vous à la connexion en double-cliquant.

Connexion au serveur de requêtes Spark en utilisant le code Java ( JDBC Client)

Assurez-vous que le CLASSPATH de Java inclut le pilote JDBC téléchargé. Exemple :

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

Vous pouvez spécifier les valeurs des paramètres et utiliser le code Java suivant pour vous connecter au serveur de requêtes Spark. Lorsque vous utilisez l'API v2, attribuez au paramètre <api_version> la valeur v2; pour l'API v3, attribuez-lui la valeur 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();
        }
    }

Connexion au serveur de requêtes Spark en utilisant Python ( PyHive JDBC Client)

Pour se connecter au serveur de requêtes Spark à l'aide d'un programme Python, procédez comme suit :

  1. Assurez-vous que vous disposez de la version Python 3.12 ou d'une version inférieure.

  2. Installer pyHive en utilisant pip install thrift "PyHive[hive_pure_sasl]==0.7.0".

  3. Sauvegardez le suivi dans un fichier tel que 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. Exécuter en utilisant python connect.py.