Améliorer la soumission d'applications Spark en utilisant l'extension de contrôle d'accès Spark

S'applique à: Moteur à étincelles

Lorsque vous soumettez une application Spark qui utilise des buckets de stockage externes enregistrés dans watsonx.data, l'extension de contrôle d'accès Spark permet une autorisation supplémentaire, ce qui renforce la sécurité. Si vous activez l'extension dans la configuration de Spark, seuls les utilisateurs autorisés peuvent accéder aux catalogues watsonx.data et les exploiter par le biais de travaux Spark.

Vous pouvez activer l'extension de contrôle d'accès Spark pour les catalogues Iceberg, Hive, Hudi et Delta Lake.

Vous pouvez utiliser Ranger ou les stratégies de données du système de gestion des accès (AMS) pour accorder ou refuser l'accès aux utilisateurs, aux groupes d'utilisateurs, aux catalogues (Iceberg, Hive, Hudi et Delta Lake ), aux schémas, aux tables et aux colonnes. Outre l'autorisation au niveau des données, le privilège de stockage est également pris en compte. Pour plus d'informations concernant l'utilisation d'AMS sur les catalogues (Iceberg, Hive, Hudi et Delta Lake ), les buckets, les schémas et les tables, voir Gestion des rôles et des privilèges. Pour plus d'informations sur la façon de créer des politiques Ranger (définies sous Hadoop SQL service) et de les activer sur les catalogues (Iceberg, Hive et Hudi), les buckets, les schémas et les tables, voir Gestion des politiques Ranger.

Prérequis

  • Créez Cloud Object Storage pour stocker les données utilisées dans l'application Spark. Pour créer Cloud Object Storage et un seau, voir Création d'un seau de stockage. Vous pouvez provisionner deux buckets, data-bucket pour stocker les tables watsonx.data et application bucket pour maintenir le code de l'application Spark.
  • Enregistrez le Cloud Object Storage bucket dans watsonx.data. Pour plus d'informations, voir Ajouter une paire de catalogues de seaux.
  • Téléchargez l'application Spark vers le stockage, voir Téléchargement des données.
  • Vous devez avoir le rôle d'administrateur IAM ou le rôle MetastoreAdmin pour créer un schéma ou une table dans watsonx.data.

Procédure

L'extension du contrôle d'accès Spark prend en charge le moteur Spark natif.

  1. Pour activer l'extension de contrôle d'accès de Spark, vous devez mettre à jour la configuration de Spark avec add authz.IBMSparkACExtension to spark.sql.extensions.

  2. Enregistrez l'application Python suivante sous iceberg.py.

L'iceberg est considéré comme un exemple. Vous pouvez également utiliser les catalogues Hive, Hudi et Delta Lake.


from pyspark.sql import SparkSession
import os

def init_spark():
    spark = SparkSession.builder \
        .appName("lh-spark-app") \
        .enableHiveSupport() \
        .getOrCreate()
    return spark

def create_database(spark):
    # Create a database in the lakehouse catalog
    spark.sql("create database if not exists lakehouse.demodb LOCATION 's3a://lakehouse-bucket/'")

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.demodb.testTable(id INTEGER, name VARCHAR(10), age INTEGER, salary DECIMAL(10, 2)) using iceberg").show()
    spark.sql("insert into lakehouse.demodb.testTable values(1,'Alan',23,3400.00),(2,'Ben',30,5500.00),(3,'Chen',35,6500.00)")
    spark.sql("select * from lakehouse.demodb.testTable").show()

def create_table_from_parquet_data(spark):
    # load parquet data into dataframe
    df = spark.read.option("header",True).parquet("file:///spark-vol/yellow_tripdata_2022-01.parquet")
    # write the dataframe into an Iceberg table
    df.writeTo("lakehouse.demodb.yellow_taxi_2022").create()
    # describe the table created
    spark.sql('describe table lakehouse.demodb.yellow_taxi_2022').show(25)
    # query the table
    spark.sql('select * from lakehouse.demodb.yellow_taxi_2022').count()

def ingest_from_csv_temp_table(spark):
    # load csv data into a dataframe
    csvDF = spark.read.option("header",True).csv("file:///spark-vol/zipcodes.csv")
    csvDF.createOrReplaceTempView("tempCSVTable")
    # load temporary table into an Iceberg table
    spark.sql('create or replace table lakehouse.demodb.zipcodes using iceberg as select * from tempCSVTable')
    # describe the table created
    spark.sql('describe table lakehouse.demodb.zipcodes').show(25)
    # query the table
    spark.sql('select * from lakehouse.demodb.zipcodes').show()

def ingest_monthly_data(spark):
    df_feb = spark.read.option("header",True).parquet("file:///spark-vol/yellow_tripdata_2022-02.parquet")
    df_march = spark.read.option("header",True).parquet("file:///spark-vol/yellow_tripdata_2022-03.parquet")
    df_april = spark.read.option("header",True).parquet("file:///spark-vol/yellow_tripdata_2022-04.parquet")
    df_may = spark.read.option("header",True).parquet("file:///spark-vol/yellow_tripdata_2022-05.parquet")
    df_june = spark.read.option("header",True).parquet("file:///spark-vol/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.demodb.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.demodb.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 => 'demodb.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.demodb.yellow_taxi_2022.files").show()
    # List all the snapshots
    # Expire earlier snapshots. Only latest one with compacted data is required
    # Again, List all the snapshots to see only 1 left
    spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.demodb.yellow_taxi_2022.snapshots").show()
    #retain only the latest one
    latest_snapshot_committed_at = spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.demodb.yellow_taxi_2022.snapshots").tail(1)[0].committed_at
    print (latest_snapshot_committed_at)
    spark.sql(f"CALL lakehouse.system.expire_snapshots(table => 'demodb.yellow_taxi_2022',older_than => TIMESTAMP '{latest_snapshot_committed_at}',retain_last => 1)").show()
    spark.sql("SELECT committed_at, snapshot_id, operation FROM lakehouse.demodb.yellow_taxi_2022.snapshots").show()
    # Removing Orphan data files
    spark.sql(f"CALL lakehouse.system.remove_orphan_files(table => 'demodb.yellow_taxi_2022')").show(truncate=False)
    # Rewriting Manifest Files
    spark.sql(f"CALL lakehouse.system.rewrite_manifests('demodb.yellow_taxi_2022')").show()

def evolve_schema(spark):
    # demonstration: Schema evolution
    # Add column fare_per_mile to the table
    spark.sql('ALTER TABLE lakehouse.demodb.yellow_taxi_2022 ADD COLUMN(fare_per_mile double)')
    # describe the table
    spark.sql('describe table lakehouse.demodb.yellow_taxi_2022').show(25)

def clean_database(spark):
    # clean-up the demo database
    spark.sql('drop table if exists lakehouse.demodb.testTable purge')
    spark.sql('drop table if exists lakehouse.demodb.zipcodes purge')
    spark.sql('drop table if exists lakehouse.demodb.yellow_taxi_2022 purge')
    spark.sql('drop database if exists lakehouse.demodb cascade')

def main():
    try:
        spark = init_spark()
        create_database(spark)
        list_databases(spark)
        basic_iceberg_table_operations(spark)
        # demonstration: Ingest parquet and csv data into a watsonx.data Iceberg table
        create_table_from_parquet_data(spark)
        ingest_from_csv_temp_table(spark)
        # load data for the month of February 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()

  1. Pour soumettre l'application Spark, spécifiez les valeurs des paramètres et exécutez la commande curl suivante. L'exemple suivant montre la commande pour soumettre l'application iceberg.py.

curl --request POST   --url https://<region>/lakehouse/api/<api_version>/spark_engines/<spark_engine_id>/applications \
   --header 'Authorization: Bearer <token>'   --header 'Content-Type: application/json'   --header 'Lhinstanceid: <instance_id>'   --data '{
  "application_details": {
  "conf": {
      "spark.hadoop.fs.s3a.bucket.<wxd-data-bucket-name>.endpoint": "<wxd-data-bucket-endpoint>",
      "spark.hadoop.fs.cos.<COS_SERVICE_NAME>.endpoint": "<COS_ENDPOINT>",
      "spark.hadoop.fs.cos.<COS_SERVICE_NAME>.secret.key": "<COS_SECRET_KEY>",
      "spark.hadoop.fs.cos.<COS_SERVICE_NAME>.access.key": "<COS_ACCESS_KEY>"
        "spark.hadoop.wxd.apiKey":"Basic xxx",
        "spark.sql.extensions":"<required-storage-support-extension>,authz.IBMSparkACExtension"

    },
    "application": "cos://<BUCKET_NAME>.<COS_SERVICE_NAME>/<python_file_name>",
  }
}

Valeurs des paramètres :

  • <token> pour obtenir le jeton d'accès à votre instance de service. Pour plus d'informations sur la génération du jeton, voir Génération d'un jeton.
  • <instance_id> iD de l'instance : L'ID de l'instance à partir de l' URL l'instance du cluster watsonx.data Exemple, crn:v1:staging:public:lakehouse:us-south:a/7bb9e380dc0c4bc284592b97d5095d3c:5b602d6a-847a-469d-bece-0a29124588c0::.
  • <wxd-data-bucket-endpoint>: Le nom d'hôte du point d'accès au réservoir de données mentionné ci-dessus. Exemple, s3.us-south.cloud-object-storage.appdomain.cloud pour un seau de stockage Cloud Object dans la région us-south.
  • <COS_SERVICE_NAME>: Fournir un nom de service de stockage d'objets Cloud.
  • <COS_ENDPOINT>: Fournit le point d'accès public. Pour plus d'informations, voir Endpoint.
  • <access_key>: Fournir l'identifiant de la clé d'accès (access_key_id). Pour plus d'informations, voir les informations d'identification.
  • <secret_key>: Fournir la clé d'accès secrète. Pour plus d'informations, voir les informations d'identification.
  • <BUCKET_NAME>: Le godet de stockage où réside le fichier d'application.
  • <python_file_name> nom du fichier de l'application Spark : Le nom du fichier de l'application Spark.
  • <api_version>:Lorsque vous utilisez l'API v2, définissez le paramètre <api_version> à v2; pour l'API v3, définissez-le à v3.

Limites :

  • L'utilisateur doit disposer d'un accès complet pour créer un schéma et une table.
  • Pour créer une politique de données, vous devez associer le catalogue au moteur Presto.
  • Si vous essayez d'afficher un schéma qui n'existe pas, le système génère un problème de nullpointer.
  • Vous pouvez activer l'extension de contrôle d'accès Spark pour les catalogues Iceberg, Hive, Hudi et Delta Lake.