Usando IBM Cloud Databases for PostgreSQL como metastore externa
Você pode usar IBM Cloud Databases for PostgreSQL para externalizar metadados fora do cluster IBM Analytics Engine Spark.
-
Crie uma instância IBM Cloud Databases for PostgreSQL. Consulte Databases for PostgreSQL.
Escolha as configurações com base em seus requisitos. Certitenha-se de escolher Rede pública & privada para a configuração do terminal. Depois de ter criado a instância e as credenciais de instância de serviço, faça uma nota do nome do banco de dados, porta, nome de usuário, senha e certificado.
-
Faça o upload do certificado Databases for PostgreSQL para um balde IBM Cloud Object Storage onde você está mantendo o seu código de aplicação.
Para acessar Databases for PostgreSQL, é necessário fornecer um certificado de cliente. Obter o certificado decodificado Base64 a partir das credenciais de serviço da instância Databases for PostgreSQL e fazer o upload do arquivo (nome que diz,
postgres.cert) a um balde Object Storage em um local específico do IBM Cloud. Mais tarde você precisará fazer o download deste certificado e disponibilizá-lo na instância IBM Analytics Engine Cargas de trabalho Spark para conexão com a metastore -
Personalize a instância IBM Analytics Engine para incluir o certificado Databases for PostgreSQL. Veja Customização baseada em Script.
Esta etapa personaliza a instância IBM Analytics Engine para fazer o certificado Databases for PostgreSQL disponível para todas as cargas de trabalho do Spark executadas contra a instância através do conjunto de bibliotecas.
-
Faça o upload do
customization_script.pyda página em Script baseado em customização para um balde IBM Cloud Object Storage. -
Execute
postgres-cert-customization-submit.jsonque usa a API REST spark-submit para customizar a instância. Observe que as referências de códigopostgres.certque você carregou para 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" } }Note que o nome do conjunto de bibliotecas
certificate_library_setdeve corresponder ao valor do parâmetro de conexão Databases for PostgreSQL metastoreae.spark.librarysetsque você especificou.
-
-
Especifique os seguintes parâmetros de conexão Databases for PostgreSQL como parte da carga útil do aplicativo Spark ou como padrões de instância. Certise-se de que você usa o terminal privado para o parâmetro
"spark.hadoop.javax.jdo.option.ConnectionURL"abaixo:"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" -
Configure o esquema de metastore Hive na instância Databases for PostgreSQL porque não há tabelas no esquema público do banco de dados Databases for PostgreSQL quando você cria a instância. Esta etapa executa o DDL relacionado ao esquema Hive para que os dados da metastore possam ser armazenados neles. Depois de executar o aplicativo Spark a seguir chamado
postgres-create-schema.py, você verá as tabelas de metadados Hive criadas contra o esquema "público" da instância.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() -
Agora execute o seguinte script chamado
postgres-parquet-table-create.pypara criar uma tabela de Parquet com metadados de IBM Cloud Object Storage no banco de dados 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() -
Execute o script PySpark chamado
postgres-parquet-table-select.pya seguir para acessar essa tabela Parquet com metadados de outra carga de trabalho do 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()