Usando o caso de uso do AWS EMR for Spark
A opção de registrar motores Spark externos no watsonx.data está obsoleta nesta versão e será removida na versão 2.3. watsonx.data O já inclui motores Spark integrados que você pode provisionar e usar diretamente, incluindo o motor Spark acelerado por Gluten e o motor Spark watsonx.data nativo.
O tópico fornece o procedimento para executar aplicativos Spark a partir do Amazon Web Services Elastic MapReduce (AWS EMR) para atingir os casos de uso do Spark IBM® watsonx.data:
- ingestão de dados
- consulta de dados
- manutenção de tabela
Pré-requisitos
-
Forneça a instância IBM® watsonx.data.
-
Crie um catálogo com o depósito do S3
-
Obter credenciais do depósito S3.
-
Configure o cluster do EMR no AWS Para obter mais informações, consulte Configurando um cluster EMR.
-
Busque as informações a seguir do IBM® watsonx.data:
- URL do MDS de watsonx.data. Para obter mais informações sobre como obter as credenciais do MDS, consulte Obtenção de credenciais do Serviço de Metadados(MDS).
- Credenciais MDS de watsonx.data. Para obter mais informações sobre como obter as credenciais do MDS, consulte Obtenção de credenciais do Serviço de Metadados(MDS).
A partir da watsonx.data versão 2.2.0, a autenticação usando
ibmlhapikeyeibmlhtokencomo nomes de usuário está obsoleta. Esses formatos serão descontinuados na 2.3.0 versão. Para garantir a compatibilidade com as próximas versões, use o novo formato:ibmlhapikey_<username>eibmlhtoken_<username>.
Visão geral
Para trabalhar com dados de origem que residem nos depósitos do AWS S3, é possível executar uma das seguintes maneiras:
- configurar a instância watsonx.data no AWS
- configurar a instância do watsonx.data baseada no IBM Cloud e incluir um catálogo baseado em depósito do AWS S3.
Os mecanismos de consulta watsonx.data podem executar consultas em dados de depósitos AWS S3. Em ambos os casos, é possível executar as operações de ingestão de dados e de manutenção de esquema baseadas em Iceberg usando o AWS EMR Spark
Sobre o caso de uso de amostra
O arquivo python de amostra (amazon-lakehouse.py) demonstra a criação de esquema (amazonschema), tabelas e dados de alimentação. Ele também suporta operações de manutenção de tabela Para obter mais informações sobre as funcionalidades na amostra, consulte Sobre o caso de uso da amostra..
Executando o caso de uso de amostra
Siga as etapas para executar o arquivo python de amostra do Spark
-
Conecte-se ao cluster AWS EMR. Para obter mais informações sobre como usar SSH para se conectar ao cluster EMR, consulte Configurando o cluster EMR.
-
Salve o arquivo python de amostra a seguir:
arquivo python de amostra do Spark
from pyspark.sql import SparkSession import os def init_spark(): spark = SparkSession.builder.appName("lh-hms-cloud")\ .enableHiveSupport().getOrCreate() return spark def create_database(spark): # Create a database in the lakehouse catalog spark.sql("create database if not exists lakehouse.amazonschema LOCATION 's3a://lakehouse-bucket-amz/'") def list_databases(spark): # list the database under lakehouse catalog spark.sql("show databases from lakehouse").show() def basic_iceberg_table_operations(spark): # demonstration: Create a basic Iceberg table, insert some data and then query table spark.sql("create table if not exists lakehouse.amazonschema.testTable(id INTEGER, name VARCHAR(10), age INTEGER, salary DECIMAL(10, 2)) using iceberg").show() spark.sql("insert into lakehouse.amazonschema.testTable values(1,'Alan',23,3400.00),(2,'Ben',30,5500.00),(3,'Chen',35,6500.00)") spark.sql("select * from lakehouse.amazonschema.testTable").show() def create_table_from_parquet_data(spark): # load parquet data into dataframce df = spark.read.option("header",True).parquet("s3a://source-bucket-amz/nyc-taxi/yellow_tripdata_2022-01.parquet") # write the dataframe into an Iceberg table df.writeTo("lakehouse.amazonschema.yellow_taxi_2022").create() # describe the table created spark.sql('describe table lakehouse.amazonschema.yellow_taxi_2022').show(25) # query the table spark.sql('select * from lakehouse.amazonschema.yellow_taxi_2022').count() def ingest_from_csv_temp_table(spark): # load csv data into a dataframe csvDF = spark.read.option("header",True).csv("s3a://source-bucket-amz/zipcodes.csv") csvDF.createOrReplaceTempView("tempCSVTable") # load temporary table into an Iceberg table spark.sql('create or replace table lakehouse.amazonschema.zipcodes using iceberg as select * from tempCSVTable') # describe the table created spark.sql('describe table lakehouse.amazonschema.zipcodes').show(25) # query the table spark.sql('select * from lakehouse.amazonschema.zipcodes').show() def ingest_monthly_data(spark): df_feb = spark.read.option("header",True).parquet("s3a://source-bucket-amz//nyc-taxi/yellow_tripdata_2022-02.parquet") df_march = spark.read.option("header",True).parquet("s3a://source-bucket-amz//nyc-taxi/yellow_tripdata_2022-03.parquet") df_april = spark.read.option("header",True).parquet("s3a://source-bucket-amz//nyc-taxi/yellow_tripdata_2022-04.parquet") df_may = spark.read.option("header",True).parquet("s3a://source-bucket-amz//nyc-taxi/yellow_tripdata_2022-05.parquet") df_june = spark.read.option("header",True).parquet("s3a://source-bucket-amz//nyc-taxi/yellow_tripdata_2022-06.parquet") df_q1_q2 = df_feb.union(df_march).union(df_april).union(df_may).union(df_june) df_q1_q2.write.insertInto("lakehouse.amazonschema.yellow_taxi_2022") def perform_table_maintenance_operations(spark): # Query the metadata files table to list underlying data files spark.sql("SELECT file_path, file_size_in_bytes FROM lakehouse.amazonschema.yellow_taxi_2022.files").show() # There are many smaller files compact them into files of 200MB each using the # `rewrite_data_files` Iceberg Spark procedure spark.sql(f"CALL lakehouse.system.rewrite_data_files(table => 'amazonschema.yellow_taxi_2022', options => map('target-file-size-bytes','209715200'))").show() # Again, query the metadata files table to list underlying data files; 6 files are compacted # to 3 files spark.sql("SELECT file_path, file_size_in_bytes FROM lakehouse.amazonschema.yellow_taxi_2022.files").show() # List all the snapshots # Expire earlier snapshots. Only latest one with comacted data is required # Again, List all the snapshots to see only 1 left spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.amazonschema.yellow_taxi_2022.snapshots").show() #retain only the latest one latest_snapshot_committed_at = spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.amazonschema.yellow_taxi_2022.snapshots").tail(1)[0].committed_at print (latest_snapshot_committed_at) spark.sql(f"CALL lakehouse.system.expire_snapshots(table => 'amazonschema.yellow_taxi_2022',older_than => TIMESTAMP '{latest_snapshot_committed_at}',retain_last => 1)").show() spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.amazonschema.yellow_taxi_2022.snapshots").show() # Removing Orphan data files spark.sql(f"CALL lakehouse.system.remove_orphan_files(table => 'amazonschema.yellow_taxi_2022')").show(truncate=False) # Rewriting Manifest Files spark.sql(f"CALL lakehouse.system.rewrite_manifests('amazonschema.yellow_taxi_2022')").show() def evolve_schema(spark): # demonstration: Schema evolution # Add column fare_per_mile to the table spark.sql('ALTER TABLE lakehouse.amazonschema.yellow_taxi_2022 ADD COLUMN(fare_per_mile double)') # describe the table spark.sql('describe table lakehouse.amazonschema.yellow_taxi_2022').show(25) def clean_database(spark): # clean-up the demo database spark.sql('drop table if exists lakehouse.amazonschema.testTable purge') spark.sql('drop table if exists lakehouse.amazonschema.zipcodes purge') spark.sql('drop table if exists lakehouse.amazonschema.yellow_taxi_2022 purge') spark.sql('drop database if exists lakehouse.amazonschema cascade') def main(): try: spark = init_spark() clean_database(spark) create_database(spark) list_databases(spark) basic_iceberg_table_operations(spark) # demonstration: Ingest parquet and csv data into a wastonx.data Iceberg table create_table_from_parquet_data(spark) ingest_from_csv_temp_table(spark) # load data for the month of Feburary to June into the table yellow_taxi_2022 created above ingest_monthly_data(spark) # demonstration: Table maintenance perform_table_maintenance_operations(spark) # demonstration: Schema evolution evolve_schema(spark) finally: # clean-up the demo database #clean_database(spark) spark.stop() if __name__ == '__main__': main() -
Execute os seguintes comandos para baixar o arquivo JAR do Serviço de Metadados do local para sua estação de trabalho. Você pode escolher os arquivos JAR com base na versão do Spark.
O arquivo JAR deve estar presente no local do
/home/hadoopem todos os nós do cluster Anote ospark.driver.extraClassPathe ospark.executor.extraClassPath.Por exemplo, se você estiver usando o Spark 3.x, inclua os seguintes arquivos JAR. Para o Spark 4.x, selecione a pasta apropriada no local.
wget https://github.com/IBM-Cloud/IBM-Analytics-Engine/raw/master/wxd-connectors/hms-connector/hive-exec-2.3.9-core.jar wget https://github.com/IBM-Cloud/IBM-Analytics-Engine/raw/master/wxd-connectors/hms-connector/hive-common-2.3.9.jar wget https://github.com/IBM-Cloud/IBM-Analytics-Engine/raw/master/wxd-connectors/hms-connector/hive-metastore-2.3.9.jar -
Configure os detalhes da conexão MDS no cluster AWS EMR para se conectar ao watsonx.data Serviço de Metadados (MDS). Um exemplo de comando para usar o spark-submit a partir de um cluster baseado em EMR 7.3.0 (Spark 3.5.1 ) é o seguinte:
Execute o comando do EMR no cluster EC2 para enviar a tarefa do spark de Amostra.
spark-submit \ --deploy-mode cluster \ --jars https://repo1.maven.org/maven2/org/apache/iceberg/iceberg-spark-runtime-3.4_2.12/1.4.0/iceberg-spark-runtime-3.4_2.12-1.4.0.jar,/usr/lib/hadoop/hadoop-aws.jar,/usr/share/aws/aws-java-sdk/aws-java-sdk-bundle*.jar,/usr/lib/hadoop-lzo/lib/* \ --conf spark.sql.catalogImplementation=hive \ --conf spark.driver.extraClassPath=/home/hadoop/hive-common-2.3.9.jar:/home/hadoop/hive-metastore-2.3.9.jar:/home/hadoop/hive-exec-2.3.9-core.jar \ --conf spark.executor.extraClassPath=/home/hadoop/hive-common-2.3.9.jar:/home/hadoop/hive-metastore-2.3.9.jar:/home/hadoop/hive-exec-2.3.9-core.jar \ --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \ --conf spark.sql.iceberg.vectorization.enabled=false \ --conf spark.sql.catalog.lakehouse=org.apache.iceberg.spark.SparkCatalog \ --conf spark.sql.catalog.lakehouse.type=hive \ --conf spark.hive.metastore.uris==<<change_endpoint>> \ --conf spark.hive.metastore.client.auth.mode=PLAIN \ --conf spark.hive.metastore.client.plain.username=ibmlhapikey \ --conf spark.hive.metastore.client.plain.password=<<change_pswd>> \ --conf spark.hive.metastore.use.SSL=true \ --conf spark.hive.metastore.truststore.type=JKS \ --conf spark.hive.metastore.truststore.path=file:///etc/pki/java/cacerts \ --conf spark.hive.metastore.truststore.password=changeit \ amazon-lakehouse.py
Valores de parâmetro:
- <<change_endpoint>> : O ponto de extremidade do URI do serviço de metadados para acessar o metastore. Para obter mais informações sobre como obter as credenciais do MDS, consulte Obtenção de credenciais do Serviço de Metadados(MDS).
- <<change_pswd>>: A senha para acessar o metastore. Para obter mais informações sobre como obter as credenciais do MDS, consulte Obtenção de credenciais do Serviço de Metadados(MDS).
Para executar o arquivo Python do Spark usando o cluster EMR 7.3.0 (Spark 3.5.1 ), baixe os jars do iceberg do local e siga o mesmo procedimento.