外部メタストアとしての IBM Cloud Databases for PostgreSQL の使用
IBM Cloud Databases for PostgreSQL を使用して、 IBM Analytics Engine Spark クラスターの外部化できます。
-
IBM Cloud Databases for PostgreSQL インスタンスを作成します。 Databases for PostgreSQL を参照してください。
要件に基づいて構成を選択します。 必ず、エンドポイント構成に 「両方のパブリック & プライベート・ネットワーク」 を選択してください。 インスタンスとサービス・インスタンス資格情報を作成した後、データベース名、ポート、ユーザー名、パスワード、および証明書をメモします。
-
アプリケーション・コードを保守する Databases for PostgreSQL 証明書を IBM Cloud Object Storage バケットにアップロードします。
Databases for PostgreSQLにアクセスするには、クライアント証明書を提供する必要があります。 Databases for PostgreSQL インスタンスのサービス資格情報から Base64 デコード証明書を取得し、ファイルをアップロードします (以下のように名前を付けます)。
postgres.cert) を特定の IBM Cloud ロケーションの Object Storage バケットに 後でこの証明書をダウンロードして、メタストアに接続するために IBM Analytics Engine インスタンス Spark ワークロードで使用できるようにする必要があります。 -
IBM Analytics Engine インスタンスをカスタマイズして、 Databases for PostgreSQL 証明書を含めます。 スクリプト・ベースのカスタマイズ を参照してください。
このステップでは、 IBM Analytics Engine インスタンスをカスタマイズして、ライブラリー・セットを介してインスタンスに対して実行されるすべての Spark ワークロードで Databases for PostgreSQL 証明書を使用できるようにします。
-
スクリプト・ベースのカスタマイズ のページから IBM Cloud Object Storage バケットに
customization_script.pyをアップロードします。 -
spark-submit REST API を使用してインスタンスをカスタマイズする
postgres-cert-customization-submit.jsonを実行します。 コードは、 IBM Cloud Object Storageにアップロードしたpostgres.certを参照していることに注意してください。{ "application_details": { "application": "/opt/ibm/customization-scripts/customize_instance_app.py", "arguments": ["{\"library_set\":{\"action\":\"add\",\"name\":\"certificate_library_set\",\"script\":{\"source\":\"py_files\",\"params\":[\"https://s3.direct.<CHANGME>.cloud-object-storage.appdomain.cloud\",\"<CHANGEME_BUCKET_NAME>\",\"postgres.cert\",\"<CHANGEME_ACCESS_KEY>\",\"<CHANGEME_SECRET_KEY>\"]}}}"], "py-files": "cos://CHANGEME_BUCKET_NAME.mycosservice/customization_script.py" } }ライブラリー・セット名
certificate_library_setは、指定した Databases for PostgreSQL メタストア接続パラメーターae.spark.librarysetsの値と一致する必要があることに注意してください。
-
-
Spark アプリケーション・ペイロードの一部として、またはインスタンスのデフォルトとして、以下の Databases for PostgreSQL メタストア接続パラメーターを指定します。 以下の
"spark.hadoop.javax.jdo.option.ConnectionURL"パラメーターにプライベート・エンドポイントを使用していることを確認してください。"spark.hadoop.javax.jdo.option.ConnectionDriverName": "org.postgresql.Driver", "spark.hadoop.javax.jdo.option.ConnectionUserName": "ibm_cloud_<CHANGEME>", "spark.hadoop.javax.jdo.option.ConnectionPassword": "<CHANGEME>", "spark.sql.catalogImplementation": "hive", "spark.hadoop.hive.metastore.schema.verification": "false", "spark.hadoop.hive.metastore.schema.verification.record.version": "false", "spark.hadoop.datanucleus.schema.autoCreateTables":"true", "spark.hadoop.javax.jdo.option.ConnectionURL": "jdbc:postgresql://<CHANGEME>.databases.appdomain.CHANGEME/ibmclouddb?sslmode=verify-ca&sslrootcert=/home/spark/shared/user-libs/certificate_library_set/custom/postgres.cert&socketTimeout=30", "ae.spark.librarysets":"certificate_library_set" -
インスタンスの作成時に Databases for PostgreSQL データベースのパブリック・スキーマに表がないため、 Databases for PostgreSQL インスタンスで Hive メタストア・スキーマをセットアップします。 このステップでは、 Hive スキーマ関連の DDL を実行して、メタストア・データをそれらの DDL に保管できるようにします。
postgres-create-schema.pyという名前の以下の Spark アプリケーションを実行すると、インスタンスの「public」スキーマに対して作成された Hive メタデータ表が表示されます。from pyspark.sql import SparkSession import time def init_spark(): spark = SparkSession.builder.appName("postgres-create-schema").getOrCreate() sc = spark.sparkContext return spark,sc def create_schema(spark,sc): tablesDF=spark.sql("SHOW TABLES") tablesDF.show() time.sleep(30) def main(): spark,sc = init_spark() create_schema(spark,sc) if __name__ == '__main__': main() -
次に、
postgres-parquet-table-create.pyという名前のスクリプトを実行して、 Databases for PostgreSQL データベースの IBM Cloud Object Storage からのメタデータを使用して Parquet テーブルを作成します。from pyspark.sql import SparkSession import time def init_spark(): spark = SparkSession.builder.appName("postgres-create-parquet-table-test").getOrCreate() sc = spark.sparkContext return spark,sc def generate_and_store_data(spark,sc): data =[("1","Romania","Bucharest","81"),("2","France","Paris","78"),("3","Lithuania","Vilnius","60"),("4","Sweden","Stockholm","58"),("5","Switzerland","Bern","51")] columns=["Ranking","Country","Capital","BroadBandSpeed"] df=spark.createDataFrame(data,columns) df.write.parquet("cos://<CHANGEME-BUCKET>.mycosservice/broadbandspeed") def create_table_from_data(spark,sc): spark.sql("CREATE TABLE MYPARQUETBBSPEED (Ranking STRING, Country STRING, Capital STRING, BroadBandSpeed STRING) STORED AS PARQUET location 'cos://CHANGEME-BUCKET.mycosservice/broadbandspeed/'") df2=spark.sql("SELECT * from MYPARQUETBBSPEED") df2.show() def main(): spark,sc = init_spark() generate_and_store_data(spark,sc) create_table_from_data(spark,sc) time.sleep(30) if __name__ == '__main__': main() -
別の Spark ワークロードからのメタデータを使用してこの Parquet 表にアクセスするには、
postgres-parquet-table-select.pyという次の PySpark スクリプトを実行します。from pyspark.sql import SparkSession import time def init_spark(): spark = SparkSession.builder.appName("postgres-select-parquet-table-test").getOrCreate() sc = spark.sparkContext return spark,sc def select_data_from_table(spark,sc): df=spark.sql("SELECT * from MYPARQUETBBSPEED") df.show() def main(): spark,sc = init_spark() select_data_from_table(spark,sc) time.sleep(60) if __name__ == '__main__': main()