Verbesserung der Spark-Anwendungseinreichung mit der Spark-Zugriffskontrollerweiterung

Wenn Sie eine Spark-Anwendung einreichen, die watsonx.data in registrierte externe Speicher-Buckets verwendet, ermöglicht die Spark-Zugriffskontroll-Erweiterung zusätzliche Autorisierungen und erhöht so die Sicherheit. Wenn Sie die Erweiterung in der Spark-Konfiguration aktivieren, dürfen nur autorisierte Benutzer über Spark-Jobs auf Kataloge zugreifen und diese watsonx.data bearbeiten.

Die Option, externe Spark-Engines in zu watsonx.data registrieren, ist veraltet und wird in Version 2.3 entfernt. enthält watsonx.data bereits integrierte Spark-Engines, die Sie direkt bereitstellen und verwenden können, darunter die Gluten-beschleunigte Spark-Engine und die native watsonx.data Spark-Engine.

Sie können die Spark-Zugriffskontroll-Erweiterung für Iceberg- Hive, Hudi- und Delta Lake-Kataloge aktivieren.

Sie können entweder Ranger- oder Access Management System (AMS)-Datenrichtlinien verwenden, um den Zugriff auf Benutzer, Benutzergruppen, Kataloge (Iceberg, Hive, Hudi und Delta Lake ), Schemata, Tabellen und Spalten zu gewähren oder zu verweigern. Neben der Berechtigung auf Datenebene wird auch die Speicherberechtigung berücksichtigt. Weitere Informationen über die Verwendung von AMS für Kataloge (Iceberg, Hive, Hudi und Delta Lake ), Buckets, Schemas und Tabellen finden Sie unter Verwaltung von Rollen und Rechten. Weitere Informationen zum Erstellen von Ranger-Richtlinien (definiert unter Hadoop SQL service) und zum Aktivieren dieser Richtlinien für Kataloge (Iceberg, Hive, Hudi und Delta Lake ), Buckets, Schemas und Tabellen finden Sie unter Verwalten von Ranger-Richtlinien.

Voraussetzungen

  • Erstellen Sie Cloud Object Storage, um die in der Spark-Anwendung verwendeten Daten zu speichern. Um Cloud Object Storage und einen Bucket zu erstellen, siehe Erstellen eines Speicher-Buckets. Sie können zwei Buckets bereitstellen: den Daten-Bucket zum Speichern von watsonx.data-Tabellen und den Anwendungs-Bucket zum Verwalten des Spark-Anwendungscodes.
  • Registrieren Sie Cloud Object Storage Bucket in watsonx.data. Weitere Informationen finden Sie unter Eimer-Katalogpaar hinzufügen.
  • Laden Sie die Spark-Anwendung in den Speicher hoch, siehe Daten hochladen.
  • Sie müssen die IAM-Administratorrolle oder die MetastoreAdmin-Rolle haben, um ein Schema oder eine Tabelle innerhalb von watsonx.data zu erstellen.

Vorgehensweise

Die Spark-Zugangskontrollerweiterung unterstützt die externe Spark-Engine.

  1. Um die Spark-Zugriffskontrollerweiterung zu aktivieren, müssen Sie die Spark-Konfiguration mit add authz.IBMSparkACExtension to spark.sql.extensions aktualisieren.

  2. Speichern Sie die folgende Python-Anwendung als iceberg.py.

Eisberg wird als Beispiel betrachtet. Sie können auch die Kataloge Delta Lake, Hudi und Hive verwenden.


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. Um die Spark-Anwendung zu übermitteln, geben Sie die Parameterwerte an und führen den folgenden curl-Befehl aus. Das folgende Beispiel zeigt den Befehl zur Übermittlung der Anwendung 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>",
  }
}

Ab watsonx.data Version ist die Authentifizierung mit ibmlhapikey 2.2.0 und ibmlhtoken als Benutzernamen veraltet. Diese Formate werden in 2.3.0 der Version auslaufen. Um die Kompatibilität mit kommenden Versionen zu gewährleisten, verwenden Sie das neue Format:ibmlhapikey_<username> und ibmlhtoken_<username>.

Parameterwerte:

  • <region>: Region, in der die Instanz bereitgestellt wird. Beispiel: us-south region.
  • <spark_engine_id>: Die eindeutige Kennung der Spark-Instanz. Informationen zum Abrufen der ID finden Sie unter "Spark-Engine-Details verwalten ".
  • <token> Um den Zugriffstoken für Ihre Serviceinstanz zu erhalten. Weitere Informationen zum Generieren des Tokens finden Sie unter "Generieren eines Tokens ".
  • <instance_id> instanz-ID: Die Instanz-ID aus der URL der watsonx.data. Beispiel, crn:v1:staging:public:lakehouse:us-south:a/7bb9e380dc0c4bc284592b97d5095d3c:5b602d6a-847a-469d-bece-0a29124588c0::.
  • <wxd-data-bucket-name>: Der Name des Daten-Buckets, der mit der Spark-Engine vom Infrastruktur-Manager verknüpft ist.
  • <wxd-data-bucket-endpoint>: Der Hostname des Endpunkts für den Zugriff auf den oben genannten Dateneimer. Beispiel: s3.us-south.cloud-object-storage.appdomain.cloud für einen Cloud Object Storage Bucket in der Region us-south.
  • <wxd-bucket-catalog-name>: Der Name des mit dem Dateneimer verbundenen Katalogs.
  • <wxd-catalog-metastore-host>: Der mit dem registrierten Bucket verbundene Metaspeicher.
  • <cos_bucket_endpoint>: Geben Sie den Wert des Metastore-Hosts an. Weitere Informationen finden Sie unter Speicherdetails.
  • <access_key>: Geben Sie die access_key_id an. Weitere Informationen finden Sie unter Speicherdetails.
  • <secret_key>: Geben Sie den secret_access_key an. Weitere Informationen finden Sie unter Speicherdetails.
  • <truststore_path>: Geben Sie den COS-Pfad an, in den das Trustore-Zertifikat hochgeladen wird. Zum Beispiel cos://di-bucket.di-test/1902xx-truststore.jks. Weitere Informationen zur Erstellung des Trustore finden Sie unter Import von selbstsignierten Zertifikaten.
  • <cas_endpoint>: Der Endpunkt des Data Access Service (DAS). Um den DAS-Endpunkt zu erhalten, siehe DAS-Endpunkt erhalten.
  • <username>: Der Benutzername für Ihre watsonx.data-Instanz. Hier, ibmlhapikey.
  • <apikey> : Die base64 kodierte `ibmlhapikey_<user_id>:<IAM_APIKEY>. Hier ist <user_id> die IBM Cloud id des Benutzers, dessen apikey für den Zugriff auf den Dateneimer verwendet wird. Um einen API-Schlüssel zu generieren, melden Sie sich in der Konsole watsonx.data an und navigieren Sie zu Profil > Profil und Einstellungen > API-Schlüssel und generieren Sie einen neuen API-Schlüssel.
  • <OBJECT_NAME> Der IBM Cloud Object Storage Name.
  • <BUCKET_NAME>: Der Speicherbereich, in dem sich die Anwendungsdatei befindet.
  • <COS_SERVICE_NAME>: Der Name des Cloud Object Storage-Dienstes.
  • <python file name>: Der Name der Spark-Anwendungsdatei.
  • <api_version> wenn Sie die API v2 verwenden, setzen Sie den Parameter <api_version> auf v2; für die API v3 setzen Sie ihn auf v3.

Einschränkungen:

  • Der Benutzer muss vollen Zugriff auf das Erstellen von Schema und Tabelle haben.
  • Um eine Datenrichtlinie zu erstellen, müssen Sie den Katalog mit der Presto-Engine verknüpfen.
  • Wenn Sie versuchen, ein Schema anzuzeigen, das nicht existiert, löst das System ein Nullpointer-Problem aus.
  • Sie können die Spark-Zugriffskontroll-Erweiterung für Iceberg- Hive, Hudi- und Delta Lake-Kataloge aktivieren.