Aprimoramento do envio de aplicativos Spark usando a extensão de controle de acesso do Spark
Quando você envia um aplicativo Spark que usa buckets de armazenamento externos registrados no watsonx.data, a extensão de controle de acesso do Spark permite autorização adicional, aumentando assim a segurança. Se você habilitar a extensão na configuração do Spark, somente usuários autorizados poderão acessar e operar watsonx.data catálogos por meio de tarefas do Spark.
A opção de registrar motores Spark externos no watsonx.data está obsoleta e será removida na versão 2.3. watsonx.data O já inclui motores Spark integrados que você pode provisionar e usar diretamente, incluindo o motor Spark acelerado por Gluten e o motor Spark watsonx.data nativo.
Você pode habilitar a extensão de controle de acesso Spark para os catálogos Iceberg, Hudi Hive e Delta Lake Catalogs.
Você pode usar as políticas de dados do Ranger ou do Access Management System (AMS) para conceder ou negar acesso a usuários, grupos de usuários, catálogos (Iceberg, Hive, Hudi e Delta Lake ), esquemas, tabelas e colunas. Além da autorização em nível de dados, o privilégio de armazenamento também é considerado. Para obter mais informações relacionadas ao uso do AMS em catálogos (Iceberg, Hive, Hudi e Delta Lake ), buckets, esquemas e tabelas, consulte Gerenciar funções e privilégios. Para obter mais informações sobre como criar políticas do Ranger (definidas em Hadoop SQL service) e ativá-las em catálogos (Iceberg, Hive, Hudi e Delta Lake ), buckets, esquemas e tabelas, consulte Gerenciar políticas do Ranger.
Pré-requisitos
- Crie Cloud Object Storage para armazenar os dados usados no aplicativo Spark. Para criar o Cloud Object Storage e um bucket, consulte Criando um bucket de armazenamento. Você pode provisionar dois compartimentos, o compartimento de dados para armazenar tabelas watsonx.data e o compartimento do aplicativo para manter o código do aplicativo Spark.
- Registre o Cloud Object Storage no watsonx.data. Para obter mais informações, consulte Adicionar par de catálogos de balde.
- Faça upload do aplicativo Spark para o armazenamento, consulte Upload de dados.
- Você deve ter a função de administrador do IAM ou a função MetastoreAdmin para criar um esquema ou uma tabela dentro de watsonx.data.
Procedimento
A extensão de controle de acesso do Spark é compatível com o mecanismo externo do Spark.
-
Para ativar a extensão de controle de acesso do Spark, você deve atualizar a configuração do Spark com
add authz.IBMSparkACExtension to spark.sql.extensions. -
Salve o seguinte aplicativo Python como iceberg.py.
O iceberg é considerado um exemplo. Você também pode usar os catálogos HiveDelta Lake Hudi e.
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 o aplicativo Spark, especifique os valores dos parâmetros e execute o seguinte comando curl. O exemplo a seguir mostra o comando para enviar o aplicativo 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 da watsonx.data versão 2.2.0, a autenticação usando ibmlhapikey e ibmlhtoken como nomes de usuário está obsoleta. Esses formatos serão descontinuados na 2.3.0 versão. Para garantir a compatibilidade com as
próximas versões, use o novo formato:ibmlhapikey_<username> e ibmlhtoken_<username>.
Valores de parâmetro:
<region>: Região onde a instância é provisionada. Exemplo, regiãous-south.<spark_engine_id>: O identificador exclusivo da instância Spark. Para obter informações sobre como recuperar a ID, consulte Gerenciando os detalhes do mecanismo do Spark.<token>: Para obter o token de acesso para sua instância de serviço. Para obter mais informações sobre como gerar o token, consulte Geração de um token.<instance_id>iD da instância: O ID da instância do URL da instância do cluster watsonx.data. Exemplo, crn:v1:staging:public:lakehouse:us-south:a/7bb9e380dc0c4bc284592b97d5095d3c:5b602d6a-847a-469d-bece-0a29124588c0::.<wxd-data-bucket-name>nome do intervalo de dados associado ao mecanismo Spark a partir do Infrastructure Manager.<wxd-data-bucket-endpoint>: o nome do host do ponto de extremidade para acessar o bucket de dados mencionado acima. Por exemplo, s3.us-south.cloud-object-storage.appdomain.cloud para um bucket de armazenamento de objetos na nuvem na região us-south.<wxd-bucket-catalog-name>: o nome do catálogo associado ao bucket de dados.<wxd-catalog-metastore-host>: o metastore associado ao bucket registrado.<cos_bucket_endpoint>: Forneça o valor do host do Metastore. Para obter mais informações, consulte detalhes de armazenamento.<access_key>: Forneça o access_key_id. Para obter mais informações, consulte detalhes de armazenamento.<secret_key>: Forneça a secret_access_key. Para obter mais informações, consulte detalhes de armazenamento.<truststore_path>: Forneça o caminho do COS no qual o certificado do trustore é carregado. Por exemplo,cos://di-bucket.di-test/1902xx-truststore.jks. Para obter mais informações sobre como gerar o trustore, consulte Importando certificados autoassinados.<cas_endpoint>: O ponto de extremidade do Serviço de Acesso a Dados (DAS). Para obter o ponto de extremidade do DAS, consulte Obtenção do ponto de extremidade do DAS.<username>: O nome de usuário da sua instância watsonx.data. Aqui, ibmlhapikey.<apikey>: O base64 codificado em `ibmlhapikey_<user_id>:<IAM_APIKEY>. Aqui, <user_id> é o IBM Cloud id do usuário cuja apikey é usada para acessar o bucket de dados. Para gerar a chave de API, faça login no console watsonx.data e navegue até Profile > Profile and Settings > API Keys e gere uma nova chave de API.<OBJECT_NAME>: O IBM Cloud Object Storage nome.<BUCKET_NAME>: o bucket de armazenamento onde reside o arquivo do aplicativo.<COS_SERVICE_NAME>: o nome do serviço de armazenamento de objetos na nuvem.<python file name>: O nome do arquivo do aplicativo Spark.<api_version>ao usar a API v2, defina o parâmetro <api_version> comov2; para a API v3, defina-o comov3.
Limitações:
- O usuário deve ter acesso total para criar o esquema e a tabela.
- Para criar a política de dados, você deve associar o catálogo ao mecanismo Presto.
- Se você tentar exibir um esquema que não existe, o sistema apresentará um problema de ponteiro nulo.
- Você pode habilitar a extensão de controle de acesso Spark para os catálogos Iceberg, Hudi Hive e Delta Lake Catalogs.