IBM Cloud Data Engine 를 메타스토어로 사용하기

IBM Cloud Data Engine 를 사용하면 Hive 메타스토어와 호환되는 카탈로그에 테이블 및 뷰의 메타데이터를 저장하고 관리할 수 있습니다.

IBM Cloud Data Engine 의 각 인스턴스에는 IBM Cloud Pak for Data 클러스터에 있는 데이터의 테이블 정의를 등록하고 관리하는 데 사용할 수 있는 데이터베이스 카탈로그가 포함되어 있습니다. 카탈로그 구문은 Hive 메타스토어 구문과 호환됩니다. IBM Cloud Data Engine 를 사용하면 IBM Cloud Pak for Data 클러스터 외부로 메타데이터를 내보낼 수 있습니다.

전제조건

Analytics Engine 에서 Data Engine을 사용하기 전에, 관련 라이브러리와 패키지를 Analytics Engine 환경으로 옮겨야 합니다.

Spark 환경에서 Data Engine에 액세스하기 위한 전제 조건

노트북을 통해 Spark 환경에서 Data Engine에 액세스하는 경우, 노트북의 셀에서 다음 명령어를 실행해야 합니다.

spark.stop()
!mkdir /tmp/dataengine
!wget https://us.sql-query.cloud.ibm.com/download/catalog/dataengine_spark-1.3.0-py3-none-any.whl -O /tmp/dataengine/dataengine_spark-1.3.0-py3-none-any.whl
# user-libs/spark2 is in the classpath of Spark
!wget https://us.sql-query.cloud.ibm.com/download/catalog/dataengine-spark-integration-1.3.0.jar -O user-libs/spark2/dataengine-spark-integration-1.3.0.jar
!wget https://us.sql-query.cloud.ibm.com/download/catalog/hive-metastore-standalone-client-3.1.2-sqlquery.jar -O  /tmp/dataengine/hive-metastore-standalone-client-3.1.2-sqlquery.jar
!pip install --force-reinstall /tmp/dataengine/dataengine_spark-1.3.0-py3-none-any.whl --user
참고: 메타스토어 관련 호출을 수행하기 전에 커널을 재시작해야 합니다.

Analytics Engine 인스턴스에서 Data Engine에 액세스하기 위한 전제 조건

애플리케이션에서 Analytics Engine 인스턴스를 통해 Data Engine에 액세스하는 경우, 해당 환경에서 Data Engine 라이브러리와 패키지를 실행할 수 있도록 이를 전송해야 합니다:

  1. 다음 명령어를 실행하여 다음 세 개의 파일을 로컬 시스템에 다운로드하십시오:

    wget https://us.sql-query.cloud.ibm.com/download/catalog/dataengine_spark-1.3.0-py3-none-any.whl
    wget https://us.sql-query.cloud.ibm.com/download/catalog/dataengine-spark-integration-1.3.0.jar
    wget https://us.sql-query.cloud.ibm.com/download/catalog/hive-metastore-standalone-client-3.1.2-sqlquery.jar
    
  2. 파일 목록의 맨 위에 있는 패키지는 “whl” 형식의 파이썬 패키지입니다. ‘cc-home’ 볼륨을 사용하여 Spark 애플리케이션 및 노트북 사용자 지정하기 문서의 단계에 따라 이 패키지를 설치하십시오.

  3. 이 파일들 중 두 번째와 세 번째 패키지는 JAR 파일이며, Spark 애플리케이션을 실행할 때 Spark 클래스 경로에 포함되어야 합니다. 파일을 사용할 수 있게 하려면:

  4. JAR 파일을 서비스 볼륨 인스턴스에 업로드하십시오.

  5. Spark 애플리케이션의 매개변수에 파일 경로를 지정하십시오:

    spark.driver.extraClassPath=</path/to/file1/in/voume:/path/to/file2/in/voume>
    spark.executor.extraClassPath=</path/to/file1/in/voume:/path/to/file2/in/voume>
    

서비스 볼륨을 생성하고 파일을 업로드하는 단계는 ‘Volumes API를 사용한 영구 볼륨 인스턴스 관리 ’에서 확인할 수 있습니다.

IBM Cloud Data Engine 를 메타스토어로 사용하기

IBM Cloud Data Engine 를 메타스토어로 사용하려면:

  1. IBM Cloud Data Engine 인스턴스를 생성합니다.IBM Cloud Data Engine 을 참조하십시오.

    IBM Cloud Data Engine 인스턴스를 프로비저닝한 후:

    1. 인스턴스의 CRN을 기록해 두십시오.
    2. 인스턴스에 대한 액세스 권한이 있는 계정 수준 API 키 또는 서비스 ID 수준 API 키를 생성합니다.
    3. 이 서비스 ID에는 IBM Cloud Data Engine 인스턴스와 버킷 모두에 대한 액세스 권한이 부여되어야 합니다.

    그런 다음 필요에 따라 인스턴스 수준이나 애플리케이션 수준에서 기본 메타스토어 구성을 사용하도록 IBM Analytics Engine powered by Apache Spark 인스턴스를 구성할 수 있습니다.

  2. IBM Cloud Data Engine 메타스토어 연결 매개변수를 지정하십시오. 다음 매개변수는 모든 애플리케이션에서 메타스토어 데이터에 액세스하려는 경우, Spark 애플리케이션 페이로드의 일부로 전달하거나 인스턴스 기본값으로 지정해야 하는 추가적인 IBM Cloud Data Engine 메타스토어 매개변수입니다:

    "spark.hive.metastore.truststore.password" : "changeit",
    "spark.hive.execution.engine":"spark",
    "spark.hive.metastore.client.plain.password":"<APIKEY-WITH-ACCESS-TO-DATA-ENGINE-INSTANCE>",
    "spark.hive.metastore.uris":"CHANGEME-thrift://catalog.us.dataengine.cloud.ibm.com:9083-CHANGEME",
    "spark.hive.metastore.client.auth.mode":"PLAIN",
    "spark.hive.metastore.use.SSL":"true",
    "spark.hive.stats.autogather":"false",
    "spark.hive.metastore.client.plain.username":"CHANGEME-crn:v1:bluemix:public:sql-query:us-south:a/abcdefgh::CHANGEME",
    # for spark 3.4 with java 11
    "spark.hive.metastore.truststore.path":"/opt/ibm/jdk/lib/security/cacerts",
    # for spark 3.4 with java 8
    "spark.hive.metastore.truststore.path":"file:///opt/ibm/jdk/jre/lib/security/cacerts",
    "spark.sql.catalogImplementation":"hive",
    "spark.sql.hive.metastore.jars":"CHANGEME --> /tmp/dataengine/*  OR path in volume",
    "spark.sql.hive.metastore.version":"3.0",
    "spark.sql.warehouse.dir":"file:///tmp",
    "spark.sql.catalogImplementation":"hive",
    "spark.hadoop.metastore.catalog.default":"spark"
    

    변수에 대해서는:

    • CHANGEME-thrift: 해당 지역의 thrift 엔드포인트를 사용하십시오. 유효한 값에 대해서는 ‘ Apache Spark 와 Data Engine 연결’을 참조하십시오.
    • CHANGEME-crn: 서비스 인스턴스 세부 정보에서 CRN을 선택하세요
  3. 데이터를 생성하고 저장합니다. 다음 정규 PySpark 애플리케이션을 실행하십시오. 이 예제에서는 generate-and-store-data.py 라는 이름으로, Parquet 데이터를 의 특정 위치에 저장합니다.

    from pyspark.sql import SparkSession
    
    def init_spark():
        spark = SparkSession.builder.appName("dataengine-generate-store-parquet-data").getOrCreate()
        sc = spark.sparkContext
        return spark,sc
    
    def generate_and_store_data(spark):
        data =[("India","New Delhi"),("France","Paris"),("Lithuania","Vilnius"),("Sweden","Stockholm"),("Switzerland","Bern")]
        columns=["Country","Capital"]
        df=spark.createDataFrame(data,columns)
        df.write.mode("overwrite").parquet("cos://mybucket.mycosservice/countriescapitals.parquet")
    
    def main():
        spark,sc = init_spark()
        generate_and_store_data(spark,sc)
    if __name__ == '__main__':
        main()
    
  4. 테이블 스키마 정의를 생성합니다. IBM Cloud Data Engine 를 메타스토어로 사용할 때는 표준 Spark SQL 구문을 사용하여 테이블을 생성할 수 없다는 점에 유의하십시오. 테이블을 만드는 방법에는 두 가지가 있습니다:

    • IBM Cloud Data Engine 사용자 인터페이스에서, 또는 표준 IBM Cloud Data Engine API( Data Engine 서비스 REST V3 API 참조)나 Python SDK( ibmcloudsql 참조)를 사용하여 수행할 수 있습니다.

      CREATE TABLE COUNTRIESCAPITALS (Country string,Capital string) 
      USING PARQUET 
      LOCATION cos://us-geo/mybucket/countriescapitals.parquet
      
    • PySpark 애플리케이션 내에서 다음 코드 스니펫을 사용하여 PySpark 에 대한 요청을 프로그래밍 방식으로 수행할 수 있습니다 create_table_data_engine.py:

      import requests
      import time
      def create_data_engine_table(api_key,crn):
          headers = {
          'Authorization': 'Basic Yng6Yng=',
          }
          data = {
          'apikey': api_key,
          'grant_type': 'urn:ibm:params:oauth:grant-type:apikey',
          }
          response = requests.post('https://iam.cloud.ibm.com/identity/token', headers=headers, data=data)
          token = response.json()['token']
      
          headers_token = {
          'Accept': 'application/json',
          'Authorization: Bearer ${TOKEN}",
          }
          params = {
          'instance_crn': crn,
          }
          json_data = {
          'statement': 'CREATE TABLE COUNTRIESCAPITALS (Country string,Capital string) USING PARQUET LOCATION cos://us-geo/mybucket/countriescapitals.parquet',
          }
          response = requests.post('https://api.dataengine.cloud.ibm.com/v3/sql_jobs', params=params, headers=headers_token, json=json_data)
          job_id = response.json()['job_id']
          time.sleep(10)
          response = requests.get(f'https://api.dataengine.cloud.ibm.com/v3/sql_jobs/{job_id}', params=params, headers=headers_token)
          if(response.json()['status']=='completed'):
              print(response.json())
      

      위치 URI(cos://us-geo/mybucket/countriescapitals.parquet)의 경우 표준 별칭 중 하나를 전달해야 한다는 점에 유의하십시오. Cloud Object Storage 엔드포인트를 참조하십시오.

      위 애플리케이션의 페이로드에는 정확한 표준 별칭(이 경우 " create_table_data_engine_payload.json us-geo")이 포함된 자격 증명을 제공해야 합니다

      {
          "application_details": {
              "conf": {
                  "spark.hadoop.fs.cos.us-geo.endpoint": "CHANGEME",
                  "spark.hadoop.fs.cos.us-geo.access.key": "CHANGEME",
                  "spark.hadoop.fs.cos.us-geo.secret.key": "CHANGEME"
              },
          "application": "cos://mybucket.us-geo/create_table_data_engine.py",
          "arguments": ["crn:v1:bluemix:public:sql-query:us-south:a/<CRN-DATA-ENGINE-INSTANCE>::","<APIKEY-WITH-ACCESS-TO-DATA-ENGINE-INSTANCE>"]
          }
      }
      
  5. 다음 애플리케이션에서 Spark SQL을 사용하여 select_query_data_engine.py테이블의 데이터를 읽어보세요:

    from pyspark.sql import SparkSession
    import time
    
    def init_spark():
      spark = SparkSession.builder.appName("dataengine-table-select-test").getOrCreate()
      sc = spark.sparkContext
      return spark,sc
    
    def select_query_data_engine(spark,sc):
      tablesDF=spark.sql("SHOW TABLES")
      tablesDF.show()
      statesDF=spark.sql("SELECT * from COUNTRIESCAPITALS");
      statesDF.show()
    
    def main():
      spark,sc = init_spark()
      select_query_data_engine(spark,sc)
    
    if __name__ == '__main__':
      main()
    

    SELECT 명령어가 정상적으로 작동하려면 식별자를 표준 별칭 중 하나로 지정해야 한다는 점에 유의하십시오. 이 예제에서는. us-geo을 사용했습니다. 예상된 조건을 충족하지 못하면 다음과 같은 오류가 표시될 수 있습니다: Configuration parse exception: Access KEY is empty. Please provide valid access key.

    select_query_data_engine_payload.json:

    {
        "application_details": {
            "conf": {
                "spark.hadoop.fs.cos.us-geo.endpoint": "CHANGEME",
                "spark.hadoop.fs.cos.us-geo.access.key": "CHANGEME",
                "spark.hadoop.fs.cos.us-geo.secret.key": "CHANGEME",
                "spark.hive.metastore.truststore.password" : "changeit",
                "spark.hive.execution.engine":"spark",
                "spark.hive.metastore.client.plain.password":"APIKEY-WITH-ACCESS-TO-DATA-ENGINE-INSTANCE",
                "spark.hive.metastore.uris":"thrift://catalog.us.dataengine.cloud.ibm.com:9083",
                "spark.hive.metastore.client.auth.mode":"PLAIN",
                "spark.hive.metastore.use.SSL":"true",
                "spark.hive.stats.autogather":"false",
                "spark.hive.metastore.client.plain.username":"crn:v1:bluemix:public:sql-query:us-south:a/<CRN-DATA-ENGINE-INSTANCE>::",
                # for spark 3.4 with java 11
                "spark.hive.metastore.truststore.path":"/opt/ibm/jdk/lib/security/cacerts",
                # for spark 3.2 with java 8
                "spark.hive.metastore.truststore.path":"file:///opt/ibm/jdk/jre/lib/security/cacerts",
                "spark.sql.catalogImplementation":"hive",
                "spark.sql.hive.metastore.jars":"/opt/ibm/connectors/data-engine/hms-client/*",
                "spark.sql.hive.metastore.version":"3.0",
                "spark.sql.warehouse.dir":"file:///tmp",
                "spark.sql.catalogImplementation":"hive",
                "spark.hadoop.metastore.catalog.default":"spark"
            },
            "application": "cos://mybucket.us-geo/select_query_data_engine.py"
        }
    }
    

편의성 API

애플리케이션에서 연결 API를 지정하여 Hive 메타스토어를 간단히 테스트해 보고 싶다면, 다음 PySpark 예제에 나와 있는 편의 API를 사용할 수 있습니다.

이 예제에서는 애플리케이션에 Hive 메타스토어 매개변수를 전달할 필요가 없습니다. 함수를 호출하면 추가적인 SparkSessionWithDataengine.enableDataengineHive 메타스토어 매개변수 없이 에 대한 연결이 초기화됩니다.

dataengine-job-convenience_api.py:

from dataengine import SparkSessionWithDataengine
from pyspark.sql import SQLContext
import sys
 from pyspark.sql import SparkSession
import time

def dataengine_table_test(spark,sc):
  tablesDF=spark.sql("SHOW TABLES")
  tablesDF.show()
  statesDF=spark.sql("SELECT * from COUNTRIESCAPITALS");
  statesDF.show()

def main():
  if __name__ == '__main__':
    if len (sys.argv) < 2:
        exit(1)
    else:
        crn = sys.argv[1]
        apikey = sys.argv[2]
        session_builder = SparkSessionWithDataengine.enableDataengine(crn, apikey, "public", "/opt/ibm/connectors/data-engine/hms-client")
        spark = session_builder.appName("Spark DataEngine integration test").getOrCreate()
        sc = spark.sparkContext
        dataengine_table_test (spark,sc)

if __name__ == '__main__':
  main()

참고로 다음은 dataengine-job-convenience_api.py페이로드입니다:

{
    "application_details": {
        "conf": {
            "spark.hadoop.fs.cos.us-geo.endpoint": "CHANGEME",
            "spark.hadoop.fs.cos.us-geo.access.key": "CHANGEME",
            "spark.hadoop.fs.cos.us-geo.secret.key": "CHANGEME"
        },
        "application": "cos://mybucket.us-geo/dataengine-job-convenience_api.py",
        "arguments": ["crn:v1:bluemix:public:sql-query:us-south:a/<CRN-DATA-ENGINE-INSTANCE>::","<APIKEY-WITH-ACCESS-TO-DATA-ENGINE-INSTANCE>"]
    }
}