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
- Instale o site watsonx.data.
- Provisione o mecanismo Spark nativo em watsonx.data.
- Faça o download do cliente JDBC:
queryserver-jdbc-SNAPSHOT-standalone.jarno link Download. - Execute o Spark Query Server no mecanismo Spark. Para criar um novo servidor de consultas, consulte Criar um servidor de consultas do Spark.
- 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.
-
Abra o DBeaver e, na barra de menus, clique em Database > Driver Manager.
-
Procurar por Hive. Você pode encontrar o driver Apache Hive 4+ na Hadoop categoria.
-
Clique em Copiar.
-
Altere o nome para Spark watsonx.data.
-
Altere as seguintes configurações :
-
Na guia Configurações,
-
Na guia Libraries (Bibliotecas )- Adicione o arquivo JAR do servidor de consulta Spark JDBC.
-
Selecione Database Navigator, clique em New Connection (Nova conexão ) e conclua as etapas a seguir:
- Selecione o driver recém-criado.
- Clique em Conectar por e selecione URL.
- Forneça o endereço JDBC URL usando o seguinte formato:
jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>. - Selecione Authentication (Autenticação ), forneça Username como seu nome de usuário e sua chave de API do IAM como a senha.
- 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:
-
Verifique se você tem a versão Python 3.12 ou inferior.
-
Instale o site pyHive usando
pip install thrift "PyHive[hive_pure_sasl]==0.7.0". -
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() -
Execute usando
python connect.py.