Mejora del envío de aplicaciones Spark mediante la extensión de control de acceso Spark
Se aplica a: Motor de chispa
Cuando envías una aplicación Spark que utiliza buckets de almacenamiento externos registrados en watsonx.data, la extensión de control de acceso de Spark permite una autorización adicional mejorando así la seguridad. Si habilita la extensión en la configuración de Spark, solo los usuarios autorizados podrán acceder a los catálogos watsonx.data y operarlos a través de los trabajos de Spark.
Puede activar la extensión de control de acceso Spark para los catálogos Iceberg, Hive, Hudi y Delta Lake.
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 y Hudi), 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 es compatible con el motor Spark nativo.
-
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. -
Guarda la siguiente aplicación Python como iceberg.py.
Iceberg se considera un ejemplo. También puede utilizar los catálogos Hive, Hudi y 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()
- 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.hadoop.wxd.apiKey":"Basic xxx",
"spark.sql.extensions":"<required-storage-support-extension>,authz.IBMSparkACExtension"
},
"application": "cos://<BUCKET_NAME>.<COS_SERVICE_NAME>/<python_file_name>",
}
}
Valores de parámetros:
<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-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.<COS_SERVICE_NAME>: Proporcione un nombre de servicio de almacenamiento de objetos en la nube.<COS_ENDPOINT>: Proporciona el punto final público. Para más información, consulte Endpoint.<access_key>: Proporcione el access_key_id. Para más información, consulte Credenciales.<secret_key>: Proporcione la clave_de_acceso_secreto. Para más información, consulte Credenciales.<BUCKET_NAME>: El bucket de almacenamiento donde reside el archivo de la aplicación.<python_file_name>nombre del archivo de la aplicación Spark.<api_version>:Cuando utilice la API v2, establezca el parámetro <api_version> env2; para la API v3, establézcalo env3.
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 activar la extensión de control de acceso Spark para los catálogos Iceberg, Hive, Hudi y Delta Lake.