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

Lorsque vous soumettez une application Spark qui utilise des compartiments de stockage externes enregistrés dans watsonx.data, l'extension de contrôle d'accès Spark permet une autorisation supplémentaire, renforçant ainsi la sécurité. Si vous activez l'extension dans la configuration Spark, seuls les utilisateurs autorisés sont autorisés à accéder aux catalogues et à watsonx.data les utiliser via les tâches Spark.

L'option permettant d'enregistrer des moteurs Spark externes dans watsonx.data est obsolète et sera supprimée dans la version 2.3. comprend watsonx.data déjà des moteurs Spark intégrés que vous pouvez provisionner et utiliser directement, notamment le moteur Spark accéléré par Gluten et le moteur Spark watsonx.data natif.

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

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 création de stratégies Ranger (définies sous Hadoop SQL service) et leur activation sur les catalogues (Iceberg, Hive, Hudi et Delta Lake ), les buckets, les schémas et les tables, voir Gestion des stratégies 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 externe.

  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 Hive catalogues Delta Lake, Hudi et.


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.sql.catalogImplementation": "hive",
      "spark.sql.iceberg.vectorization.enabled":"false",
        "spark.sql.catalog.<wxd-bucket-catalog-name>":"org.apache.iceberg.spark.SparkCatalog",
      "spark.sql.catalog.<wxd-bucket-catalog-name>.type":"hive",
      "spark.sql.catalog.<wxd-bucket-catalog-name>.uri":"thrift://<wxd-catalog-metastore-host>",
      "spark.hive.metastore.client.auth.mode":"PLAIN",
      "spark.hive.metastore.client.plain.username":"<username>",
      "spark.hive.metastore.client.plain.password":"xxx",
      "spark.hive.metastore.use.SSL":"true",
      "spark.hive.metastore.truststore.type":"JKS",
      "spark.hive.metastore.truststore.path":"<truststore_path>",
      "spark.hive.metastore.truststore.password":"changeit",
        "spark.hadoop.fs.s3a.bucket.<wxd-data-bucket-name>.aws.credentials.provider":"com.ibm.iae.s3.credentialprovider.WatsonxCredentialsProvider",
        "spark.hadoop.fs.s3a.bucket.<wxd-data-bucket-name>.custom.signers":"WatsonxAWSV4Signer:com.ibm.iae.s3.credentialprovider.WatsonxAWSV4Signer",
        "spark.hadoop.fs.s3a.bucket.<wxd-data-bucket-name>.s3.signing-algorithm":"WatsonxAWSV4Signer",
        "spark.hadoop.wxd.cas.endpoint":"<cas_endpoint>/cas/v1/signature",
        "spark.hadoop.wxd.instanceId":"<instance_crn>",
        "spark.hadoop.wxd.apiKey":"Basic xxx",
        "spark.wxd.api.endpoint":"<wxd-endpoint>",
        "spark.driver.extraClassPath":"opt/ibm/connectors/wxd/spark-authz/cpg-client-1.0-jar-with-dependencies.jar:/opt/ibm/connectors/wxd/spark-authz/ibmsparkacextension_2.12-1.0.jar",
        "spark.sql.extensions":"<required-storage-support-extension>,authz.IBMSparkACExtension"

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

À partir de watsonx.data la version 2.2.0, l'authentification à l'aide de ibmlhapikey et ibmlhtoken comme noms d'utilisateur est obsolète. Ces formats sont progressivement supprimés dans 2.3.0 la version. Pour assurer la compatibilité avec les prochaines versions, utilisez le nouveau format :ibmlhapikey_<username> et ibmlhtoken_<username>.

Valeurs des paramètres :

  • <region>: Région dans laquelle l'instance est provisionnée. Exemple, région d' us-south.
  • <spark_engine_id>: Identifiant unique de l'instance Spark. Pour savoir comment récupérer l'ID, consultez la rubrique Gérer les détails du moteur Spark.
  • <token>: Pour obtenir le jeton d'accès pour 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-name>: Nom du compartiment de données associé au moteur Spark par le gestionnaire d'infrastructure.
  • <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.
  • <wxd-bucket-catalog-name>: Le nom du catalogue associé au seau de données.
  • <wxd-catalog-metastore-host>: Le métastore associé au seau enregistré.
  • <cos_bucket_endpoint>: Fournir la valeur de l'hôte du Metastore. Pour plus d'informations, voir détails du stockage.
  • <access_key>: Fournir l'identifiant de la clé d'accès (access_key_id). Pour plus d'informations, voir détails du stockage.
  • <secret_key>: Fournir la clé d'accès secrète. Pour plus d'informations, voir détails du stockage.
  • <truststore_path>: Indiquer le chemin COS où le certificat trustore est téléchargé. Par exemple cos://di-bucket.di-test/1902xx-truststore.jks. Pour plus d'informations sur la génération du trustore, voir Importer des certificats auto-signés.
  • <cas_endpoint>: Le point d'accès au service d'accès aux données (DAS). Pour obtenir le point de terminaison DAS, voir Obtention du point de terminaison DAS.
  • <username>: Le nom d'utilisateur de votre instance watsonx.data. Ici, ibmlhapikey.
  • <apikey> : Le base64 encodé `ibmlhapikey_<user_id>:<IAM_APIKEY>. Ici, <user_id> est l'identifiant IBM Cloud de l'utilisateur dont l'apikey est utilisé pour accéder au seau de données. Pour générer une clé API, connectez-vous à la console watsonx.data et naviguez vers Profil > Profil et paramètres > Clés API et générez une nouvelle clé API.
  • <OBJECT_NAME>: Le IBM Cloud Object Storage nom.
  • <BUCKET_NAME>: Le godet de stockage où réside le fichier d'application.
  • <COS_SERVICE_NAME>: Le nom du service de stockage d'objets dans le nuage.
  • <python file name>: Nom du fichier d'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 Delta Lake Iceberg, Hive Hudi et.