Spark JDBC 드라이버를 사용하여 Spark 쿼리 서버에 연결하기

에 적용됩니다: 스파크 엔진

다음과 같은 방법으로 Spark 쿼리 서버에 연결하고 쿼리를 실행하여 데이터를 분석할 수 있습니다.

시작하기 전에

  1. 설치 watsonx.data.
  2. watsonx.data 에서 네이티브 Spark 엔진을 프로비저닝합니다.
  3. 다운로드 링크에서 JDBC 클라이언트( queryserver-jdbc-SNAPSHOT-standalone.jar )를 다운로드하세요.
  4. 스파크 엔진에서 스파크 쿼리 서버를 실행합니다. 새 쿼리 서버를 만들려면 Spark 쿼리 서버 만들기를 참조하세요.
  5. 연결 속성 - 쿼리 서버의 점 3개 메뉴를 클릭하고 연결 세부정보를 클릭한 다음 다음 연결 세부정보를 복사합니다:
    • host
    • URI
    • 인스턴스
    • 사용자 이름
    • IBM IAM API 키입니다.

DBeaver( JDBC 클라이언트)를 사용하여 Spark 쿼리 서버에 연결하기

DBeaver와 같은 JDBC 클라이언트를 사용하여 Spark 쿼리 서버에 연결하려면 DBeaver에서 watsonx.data 드라이버를 설정합니다.

  1. DBeaver를 열고 메뉴 모음에서 데이터베이스 > 드라이버 관리자를 클릭합니다.

  2. 검색 Hive. Apache Hive 4+ 드라이버는 Hadoop 카테고리에서 찾을 수 있습니다.

  3. 복사를 클릭하십시오.

  4. 이름을 Spark watsonx.data 로 변경합니다.

  5. 다음 설정을 변경합니다 :

  6. 설정 탭에서

  7. 라이브러리 탭에서 - Spark JDBC 쿼리 서버 JAR 파일을 추가합니다.

  8. 데이터베이스 Navigator 를 선택하고 새 연결을 클릭한 후 다음 단계를 완료합니다:

    1. 새로 생성된 드라이버를 선택합니다.
    2. 연결 기준 을 클릭하고 URL.
    3. 다음 형식을 사용하여 JDBC URL 을 제공하십시오: jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>.
    4. 인증을 선택하고 사용자 아이디를 사용자 이름으로, IAM API 키를 비밀번호로 입력합니다.
    5. 저장하고 두 번 클릭하여 연결에 연결합니다.

Java ( JDBC 클라이언트) 코드를 사용하여 Spark 쿼리 서버에 연결하기

Java 클래스 경로에 다운로드한 JDBC 드라이버가 포함되어 있는지 확인합니다. 예를 들어, 다음과 같습니다.

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

매개변수 값을 지정하고 다음 Java 코드를 사용하여 Spark 쿼리 서버에 연결할 수 있습니다. v2 API를 사용하는 경우 <api_version> 파라미터를 v2 으로 설정하고, v3 API의 경우 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();
        }
    }

Python ( PyHive JDBC 클라이언트)를 사용하여 Spark 쿼리 서버에 연결합니다

Python 프로그램을 사용하여 Spark 쿼리 서버에 연결하려면 다음과 같이 하세요:

  1. Python 버전 3.12 이하가 있는지 확인합니다.

  2. pip install thrift "PyHive[hive_pure_sasl]==0.7.0" 을 사용하여 pyHive 을 설치합니다.

  3. 팔로우를 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. python connect.py 을 사용하여 실행합니다.