Spark JDBC Driver を使用して Spark クエリサーバーに接続する

適用対象スパークエンジン

以下の方法でSparkクエリサーバーに接続し、クエリを実行してデータを分析できます。

開始前に

  1. watsonx.data をインストールする。
  2. watsonx.data でネイティブ・スパーク・エンジンを提供する。
  3. ダウンロード リンクから JDBC クライアント: queryserver-jdbc-SNAPSHOT-standalone.jar をダウンロードしてください。
  4. SparkエンジンでSpark Query Serverを実行します。 新しいクエリサーバーを作成するには、 Sparkクエリサーバーを作成するを 参照してください。
  5. 接続プロパティ - Query Serverの3点メニューをクリックし、[ 接続の詳細] をクリックして、以下の接続の詳細をコピーします:
    • ホスト
    • 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. Libraries タブで、Spark JDBC query server JAR ファイルを追加します。

  8. Database Navigator を選択し、 New Connection をクリックし、以下のステップを完了する:

    1. 新しく作成したドライバーを選択します。
    2. で接続をクリックし、以下を選択します。 URL.
    3. JDBC URL 以下の書式で記入: jdbc:hive2://<HOST>:443/default;instance=<INSTANCE>;httpPath=<URI>.
    4. Authenticationを選択し、ユーザー名に Username、パスワードにIAM APIキーを入力します。
    5. 保存し、ダブルクリックで接続する。

Java ( JDBC Client) コードを使用して Spark クエリサーバーに接続する

Java CLASSPATHにダウンロードした 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 Client ) を使用して 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 を使って実行する。