Utilizzo di AWS EMR per il caso di utilizzo Spark
L'opzione di registrare motori Spark esterni in watsonx.data è deprecata in questa versione e verrà rimossa nella versione 2.3. include watsonx.data già motori Spark integrati che è possibile fornire e utilizzare direttamente, tra cui il motore Spark accelerato da Gluten e il motore Spark watsonx.data nativo.
L'argomento fornisce la procedura per eseguire applicazioni Spark da Amazon Web Services Elastic MapReduce (AWS EMR) per ottenere i casi di utilizzo Spark IBM® watsonx.data:
- immissione di dati
- query di dati
- manutenzione tabella
Prerequisiti
-
Esegui provisioning dell'istanza IBM® watsonx.data.
-
Crea un catalogo con il bucket S3.
-
Ottieni credenziali bucket S3.
-
Configurare il cluster EMR su AWS. Per ulteriori informazioni, consultare Impostazione di un cluster EMR.
-
Recupera le seguenti informazioni da IBM® watsonx.data:
- URL MDS da watsonx.data. Per ulteriori informazioni su come ottenere le credenziali MDS, vedere Ottenere le credenziali di Metadata Service(MDS).
- Credenziali MDS da watsonx.data. Per ulteriori informazioni su come ottenere le credenziali MDS, vedere Ottenere le credenziali di Metadata Service(MDS).
A partire dalla watsonx.data versione 2.2.0, l'autenticazione utilizzando
ibmlhapikeyeibmlhtokencome nomi utente è deprecata. Questi formati vengono gradualmente eliminati nella 2.3.0 versione. Per garantire la compatibilità con le prossime versioni, utilizzare il nuovo formato:ibmlhapikey_<username>eibmlhtoken_<username>.
Panoramica
Per gestire i dati di origine che risiedono nei bucket AWS S3, è possibile effettuare una delle seguenti operazioni:
- configurare l'istanza watsonx.data su AWS
- configurare l'istanza IBM Cloud basata su watsonx.data e includere un catalogo basato su bucket AWS S3.
I motori di query watsonx.data possono eseguire query sui dati dai bucket AWS S3. In entrambi i casi, puoi eseguire l'inserimento dei dati e le operazioni di manutenzione dello schema basate su Iceberg utilizzando AWS EMR Spark.
Informazioni sul caso di utilizzo di esempio
Il file python di esempio (amazon-lakehouse.py) illustra la creazione dello schema (amazonschema), le tabelle e l'inserimento dei dati. Supporta anche le operazioni di manutenzione delle tabelle. Per ulteriori informazioni sulle funzionalità nell'esempio, consultare Informazioni sul caso di utilizzo di esempio
Esecuzione del caso di utilizzo di esempio
Seguire la procedura per eseguire il file python di esempio Spark.
-
Connettersi al cluster EMR AWS. Per ulteriori informazioni sull'utilizzo di SSH per la connessione al cluster EMR, vedi Impostazione del cluster EMR.
-
Salvare il seguente file python di esempio.
File python di esempio 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() -
Eseguire i seguenti comandi per scaricare il file JAR del servizio metadati dalla posizione alla workstation. È possibile scegliere i file JAR in base alla versione di Spark.
Il file JAR deve trovarsi nell'ubicazione
/home/hadoopsu tutti i nodi del cluster. Prendere nota dispark.driver.extraClassPathespark.executor.extraClassPath.Ad esempio, se si utilizza Spark 3.x, includere i seguenti file JAR. Per Spark 4.x, seleziona la cartella appropriata dalla posizione.
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 -
Configurare i dettagli della connessione MDS nel cluster AWS EMR per collegarsi al watsonx.data Metadata Service (MDS). Un comando di esempio per utilizzare spark-submit da un cluster basato su EMR 7.3.0 (Spark 3.5.1 ) è il seguente:
Eseguire il comando da EMR sul cluster EC2 per inoltrare il lavoro spark di esempio.
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
Valori di parametro:
- <<change_endpoint>> : L'endpoint URI del servizio di metadati per accedere al metastore. Per ulteriori informazioni sull'ottenimento delle credenziali MDS, vedere Ottenere le credenziali di Metadata Service(MDS).
- <<change_pswd>> : La password per accedere al metastore. Per ulteriori informazioni sull'ottenimento delle credenziali MDS, vedere Ottenere le credenziali di Metadata Service(MDS).
Per eseguire il file Python Spark utilizzando il cluster EMR 7.3.0 (Spark 3.5.1 ), scaricare i file jar Iceberg dalla posizione indicata e seguire la stessa procedura.