Conectando-se ao servidor de consulta do Spark usando o Spark JDBC Driver

Aplica-se a: Motor de faísca

Você pode se conectar ao servidor de consultas do Spark das seguintes maneiras e executar consultas para analisar seus dados.

Antes de Iniciar

  1. Instale o site watsonx.data.
  2. Provisione o mecanismo Spark nativo em watsonx.data.
  3. Faça o download do cliente JDBC: queryserver-jdbc-SNAPSHOT-standalone.jar no link Download.
  4. Execute o Spark Query Server no mecanismo Spark. Para criar um novo servidor de consultas, consulte Criar um servidor de consultas do Spark.
  5. Propriedades da conexão - Clique no menu de três pontos do Query Server, clique em Detalhes da conexão e copie os seguintes detalhes da conexão:
    • Host
    • URI
    • Instância
    • Nome do usuário
    • Sua chave de API do IAM IBM.

Conexão com o servidor de consulta Spark usando o DBeaver (cliente JDBC )

Para se conectar ao servidor de consulta do Spark usando um cliente JDBC, como o DBeaver, configure o driver watsonx.data no DBeaver.

  1. Abra o DBeaver e, na barra de menus, clique em Database > Driver Manager.

  2. Procurar por Hive. Você pode encontrar o driver Apache Hive 4+ na Hadoop categoria.

  3. Clique em Copiar.

  4. Altere o nome para Spark watsonx.data.

  5. Altere as seguintes configurações :

  6. Na guia Configurações,

  7. Na guia Libraries (Bibliotecas )- Adicione o arquivo JAR do servidor de consulta Spark JDBC.

  8. Selecione Database Navigator, clique em New Connection (Nova conexão ) e conclua as etapas a seguir:

    1. Selecione o driver recém-criado.
    2. Clique em Conectar por e selecione URL.
    3. Forneça o endereço JDBC URL usando o seguinte formato: jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>.
    4. Selecione Authentication (Autenticação ), forneça Username como seu nome de usuário e sua chave de API do IAM como a senha.
    5. Salve e conecte-se à conexão clicando duas vezes.

Conexão com o servidor de consulta Spark usando o código Java (cliente JDBC )

Verifique se o CLASSPATH do site Java inclui o driver JDBC baixado. Por exemplo:

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

Você pode especificar os valores dos parâmetros e usar o seguinte código Java para se conectar ao servidor de consulta do Spark. Ao usar a API v2, defina o parâmetro <api_version> como v2; para a API v3, defina-o como 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();
        }
    }

Conectando-se ao servidor de consulta Spark usando Python ( PyHive JDBC Client)

Para se conectar ao servidor de consulta do Spark usando um programa Python, faça o seguinte:

  1. Verifique se você tem a versão Python 3.12 ou inferior.

  2. Instale o site pyHive usando pip install thrift "PyHive[hive_pure_sasl]==0.7.0".

  3. Salve o seguinte em um arquivo como 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. Execute usando python connect.py.