Utilisation de IBM Cloud Databases for PostgreSQL en tant que métamagasin externe
Vous pouvez utiliser IBM Cloud Databases for PostgreSQL pour externaliser des métadonnées en dehors du cluster Spark IBM Analytics Engine.
-
Créez une instance IBM Cloud Databases for PostgreSQL. Voir Databases for PostgreSQL.
Choisissez les configurations en fonction de vos besoins. Veillez à choisir Réseau public et réseau privé pour la configuration de noeud final. Après avoir créé l'instance et les données d'identification de l'instance de service, notez le nom de la base de données, le port, le nom d'utilisateur, le mot de passe et le certificat.
-
Téléchargez le certificat Databases for PostgreSQL dans un compartiment IBM Cloud Object Storage où vous conservez votre code d'application.
Pour accéder à Databases for PostgreSQL, vous devez fournir un certificat client. Obtenez le certificat décodé Base64 à partir des données d'identification du service de l'instance Databases for PostgreSQL et téléchargez le fichier (par exemple,
postgres.cert) à un compartiment Object Storage dans un emplacement IBM Cloud spécifique. Vous devrez ensuite télécharger ce certificat et le rendre disponible dans les charges de travail Spark de l'instance IBM Analytics Engine pour la connexion au métamagasin -
Personnalisez l'instance IBM Analytics Engine pour inclure le certificat Databases for PostgreSQL. Voir Personnalisation basée sur un script.
Cette étape personnalise l'instance IBM Analytics Engine pour rendre le certificat Databases for PostgreSQL disponible pour toutes les charges de travail Spark exécutées sur l'instance via l'ensemble de bibliothèques.
-
Téléchargez le fichier
customization_script.pyde la page Personnalisation basée sur un script dans un compartiment IBM Cloud Object Storage. -
Exécutez
postgres-cert-customization-submit.jsonqui utilise l'API REST spark-submit pour personnaliser l'instance. Notez que le code fait référence àpostgres.certque vous avez téléchargé dans IBM Cloud Object Storage.{ "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" } }Notez que le nom de l'ensemble de bibliothèques
certificate_library_setdoit correspondre à la valeur du paramètre de connexion au métamagasin Databases for PostgreSQLae.spark.librarysetsque vous avez spécifié.
-
-
Spécifiez les paramètres de connexion de métamagasin Databases for PostgreSQL suivants dans le contenu de l'application Spark ou comme paramètres par défaut de l'instance. Veillez à utiliser le noeud final privé pour le paramètre
"spark.hadoop.javax.jdo.option.ConnectionURL"ci-dessous:"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" -
Configurez le schéma de métamagasin Hive dans l'instance Databases for PostgreSQL car il n'existe aucune table dans le schéma public de la base de données Databases for PostgreSQL lorsque vous créez l'instance. Cette étape exécute le DDL lié au schéma Hive afin que les données de métamagasin puissent y être stockées. Après avoir exécuté l'application Spark suivante appelée
postgres-create-schema.py, vous verrez les tables de métadonnées Hive créées par rapport au schéma "public" de l'instance.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() -
Exécutez à présent le script suivant appelé
postgres-parquet-table-create.pypour créer une table Parquet avec les métadonnées de IBM Cloud Object Storage dans la base de données Databases for PostgreSQL.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() -
Exécutez le script PySpark suivant appelé
postgres-parquet-table-select.pypour accéder à cette table Parquet avec les métadonnées d'une autre charge de travail Spark: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()