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.

  1. 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.

  2. 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

  3. 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.

    1. Faça o upload do customization_script.py da página em Script baseado em customização para um balde IBM Cloud Object Storage.

    2. Execute postgres-cert-customization-submit.json que usa a API REST spark-submit para customizar a instância. Observe que as referências de código postgres.cert que 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_set deve corresponder ao valor do parâmetro de conexão Databases for PostgreSQL metastore ae.spark.librarysets que você especificou.

  4. 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"
    
  5. 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()
    
  6. Agora execute o seguinte script chamado postgres-parquet-table-create.py para 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()
    
  7. Execute o script PySpark chamado postgres-parquet-table-select.py a 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()