Mejora del envío de aplicaciones Spark mediante la extensión de control de acceso Spark

Cuando envías una aplicación Spark que utiliza depósitos de almacenamiento externos registrados en watsonx.data, la extensión de control de acceso de Spark permite una autorización adicional, lo que mejora la seguridad. Si habilita la extensión en la configuración de Spark, solo los usuarios autorizados podrán acceder y operar watsonx.data catálogos a través de trabajos de Spark.

La opción de registrar motores Spark externos en watsonx.data ha quedado obsoleta y se eliminará en la versión 2.3. watsonx.data ya incluye motores Spark integrados que puede aprovisionar y utilizar directamente, incluidos el motor Spark acelerado por Gluten y el motor Spark watsonx.data nativo.

Puede habilitar la extensión de control de acceso Spark para los catálogos Delta Lake Iceberg, Hive Hudi y.

Puede utilizar las políticas de datos de Ranger o de Access Management System (AMS) para conceder o denegar el acceso a usuarios, grupos de usuarios, catálogos (Iceberg, Hive, Hudi y Delta Lake ), esquemas, tablas y columnas. Además de la autorización a nivel de datos, también se tiene en cuenta el privilegio de almacenamiento. Para más información sobre el uso de AMS en catálogos (Iceberg, Hive, Hudi y Delta Lake ), buckets, esquemas y tablas, consulte Gestión de roles y privilegios. Para obtener más información sobre cómo crear políticas de Ranger (definidas en el servicio Hadoop SQL) y cómo habilitarlas en catálogos (Iceberg, Hive, Hudi y Delta Lake ), buckets, esquemas y tablas, consulte Administración de políticas de Ranger.

Requisitos previos

  • Crear Cloud Object Storage para almacenar los datos utilizados en la aplicación Spark. Para crear Cloud Object Storage y un bucket, consulta Crear un bucket de almacenamiento. Puede aprovisionar dos buckets, data-bucket para almacenar watsonx.data tablas y application bucket para mantener el código de la aplicación Spark.
  • Registrar Cloud Object Storage bucket en watsonx.data. Para obtener más información, consulte Añadir par de catálogos de cubos.
  • Cargue la aplicación Spark en el almacenamiento, consulte Carga de datos.
  • Debe tener el rol de administrador IAM o el rol MetastoreAdmin, para crear esquema o tabla dentro de watsonx.data.

Procedimiento

La extensión de control de acceso Spark admite un motor Spark externo.

  1. Para habilitar la extensión de control de acceso de Spark, debes actualizar la configuración de Spark con add authz.IBMSparkACExtension to spark.sql.extensions.

  2. Guarda la siguiente aplicación Python como iceberg.py.

Iceberg se considera un ejemplo. También puede utilizar los Hive catálogos Delta Lake, Hudi y.


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. Para enviar la aplicación Spark, especifique los valores de los parámetros y ejecute el siguiente comando curl. El siguiente ejemplo muestra el comando para enviar la aplicación 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>",
  }
}

A partir de watsonx.data la versión 2.2.0, la autenticación utilizando ibmlhapikey y ibmlhtoken como nombres de usuario queda obsoleta. Estos formatos se eliminan gradualmente en 2.3.0 la versión. Para garantizar la compatibilidad con las próximas versiones, utilice el nuevo formato:ibmlhapikey_<username> y ibmlhtoken_<username>.

Valores de parámetros:

  • <region>: Región en la que se aprovisiona la instancia. Por ejemplo, la región de us-south.
  • <spark_engine_id>: El identificador único de la instancia Spark. Para obtener información sobre cómo recuperar el ID, consulte Administrar los detalles del motor Spark.
  • <token>: Para obtener el token de acceso para su instancia de servicio. Para obtener más información sobre cómo generar el token, consulte Generar un token.
  • <instance_id> el ID de instancia de la URL de instancia del clúster watsonx.data. Ejemplo, crn:v1:staging:public:lakehouse:us-south:a/7bb9e380dc0c4bc284592b97d5095d3c:5b602d6a-847a-469d-bece-0a29124588c0::.
  • <wxd-data-bucket-name>: El nombre del depósito de datos asociado al motor de chispa del administrador de infraestructura.
  • <wxd-data-bucket-endpoint>: El nombre de host del endpoint para acceder al bucket de datos mencionado anteriormente. Ejemplo, s3.us-south.cloud-object-storage.appdomain.cloud para un bucket de almacenamiento Cloud Object en la región us-south.
  • <wxd-bucket-catalog-name>: El nombre del catálogo asociado al cubo de datos.
  • <wxd-catalog-metastore-host>: El metastore asociado al bucket registrado.
  • <cos_bucket_endpoint>: Proporcione el valor del host de Metastore. Para obtener más información, consulte detalles de almacenamiento.
  • <access_key>: Proporcione el access_key_id. Para obtener más información, consulte detalles de almacenamiento.
  • <secret_key>: Proporcione la clave_de_acceso_secreto. Para obtener más información, consulte detalles de almacenamiento.
  • <truststore_path>: Proporcione la ruta COS donde se carga el certificado trustore. Por ejemplo cos://di-bucket.di-test/1902xx-truststore.jks. Para obtener más información sobre la generación del trustore, consulte Importación de certificados autofirmados.
  • <cas_endpoint>: El punto final del servicio de acceso a datos (DAS). Para obtener el punto final DAS, consulte Obtención del punto final DAS.
  • <username>: El nombre de usuario de su instancia watsonx.data. Aquí, ibmlhapikey.
  • <apikey> : El base64 codificado `ibmlhapikey_<user_id>:<IAM_APIKEY>. Aquí, <user_id> es el IBM Cloud id del usuario cuya apikey se utiliza para acceder al bucket de datos. Para generar la clave de API, inicie sesión en la consola watsonx.data y vaya a Perfil > Perfil y configuración > Claves de API y genere una nueva clave de API.
  • <OBJECT_NAME>: El IBM Cloud Object Storage nombre.
  • <BUCKET_NAME>: El bucket de almacenamiento donde reside el archivo de la aplicación.
  • <COS_SERVICE_NAME>: El nombre del servicio de Almacenamiento de objetos en la nube.
  • <python file name>: El nombre del archivo de la aplicación Spark.
  • <api_version>:Cuando utilice la API v2, establezca el parámetro <api_version> en v2; para la API v3, establézcalo en v3.

Limitaciones:

  • El usuario debe tener acceso total para crear el esquema y la tabla.
  • Para crear la política de datos, debe asociar el catálogo al motor Presto.
  • Si se intenta mostrar un esquema que no existe, el sistema genera un problema de puntero nulo.
  • Puede habilitar la extensión de control de acceso Spark para los catálogos Delta Lake Iceberg, Hive Hudi y.