Monitoraggio dell'esecuzione delle applicazioni Spark con Databand
Si applica a: Motore a scintilla
L'integrazione di Databand con Spark migliora le capacità di monitoraggio, fornendo approfondimenti che vanno oltre l' Spark UI e la cronologia di Spark.
Databand migliora il monitoraggio delle applicazioni Spark come segue:
- Monitoraggio avanzato: Le annotazioni dei task di Databand consentono di etichettare e tracciare le fasi cruciali dell'applicazione Spark, offrendo un livello di monitoraggio più significativo rispetto ai job, agli stage o ai task di Spark.
- Tracciamento dei dataset: Databand monitora i set di dati a cui si accede e che vengono modificati durante l'esecuzione dell'applicazione Spark, fornendo una maggiore visibilità sui flussi di dati.
- Avvisi personalizzati: è possibile configurare avvisi per fasi specifiche dell'applicazione o tracciare metriche chiave del set di dati, consentendo di identificare e risolvere tempestivamente potenziali problemi.
Per iniziare a utilizzare Databand, è necessario avere un abbonamento attivo a Databand. Si può ottenere richiedendo un'istanza dell'applicazione Databand nel cloud (SaaS), che viene distribuita dal team Databand, oppure optando per un'installazione self-hosted (on-premises). Per integrare databand con la propria istanza watsonx.data, è necessario disporre delle seguenti credenziali:
- Indirizzo dell'ambiente: L' URL dell'ambiente Databand (ad esempio: yourcompanyname.databand.ai ).
- Take di accesso: Un token di accesso a Databand necessario per connettersi all'ambiente. È possibile generare e gestire i token attraverso l'interfaccia utente di Databand, come richiesto. Per maggiori dettagli, visitate il sito: Gestione dei token di accesso personali.
Procedura
- Accedi alla console watsonx.data.
- Dal menu di navigazione, andare su Configurazioni > IBM Data Observability by Databand.
- Inserire l'indirizzo dell'ambiente nel campo URL e il token di accesso nel campo Token di accesso.
- Fare clic su Prova connessione per convalidare la connessione.
- Fare clic su Salva per salvare i dettagli. È possibile modificare i dettagli facendo clic su Modifica.
Databand entra in vigore per i nuovi lavori Spark avviati dopo l'abilitazione. I lavori precedenti e in corso non vengono registrati per l'analisi.
Ascoltatore di banda dati
Questo metodo tiene automaticamente traccia delle operazioni sul set di dati. Il vostro script Spark può trarre vantaggio dal tracciamento automatico delle operazioni sui dataset.
Decoratori di database e API di registrazione
Per utilizzare questo metodo, è necessario importare il modulo dbnd, che richiede modifiche al codice.
Esempio di utilizzo end-to-end delle API Databand
Il seguente esempio dimostra l'uso delle API di dbnd.
example.py*
import time
import logging
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
from dbnd import dbnd_tracking, task, dataset_op_logger, log_metric, log_dataframe
# Initialize Spark session
spark = SparkSession.builder \
.appName("Data Pipeline with Databand") \
.getOrCreate()
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
@task
def create_sample_data():
# Create a DataFrame with sample data including columns to be dropped
data = [
("John", "Camping Equipment", 500, "Regular", "USA"),
("Jane", "Golf Equipment", 300, "Premium", "UK"),
("Mike", "Camping Equipment", 450, "Regular", "USA"),
("Emily", "Golf Equipment", 350, "Premium", "Canada"),
("Anna", "Camping Equipment", 600, "Regular", "USA"),
("Tom", "Golf Equipment", 200, "Regular", "UK")
]
columns = ["Name", "Product line", "Sales", "Customer Type", "Country"]
retailData = spark.createDataFrame(data, columns)
# Log the data creation
unique_file_name = "sample-data"
with dataset_op_logger(unique_file_name, "read", with_schema=True, with_preview=True, with_stats=True) as logger:
logger.set(data=retailData)
return retailData
@task
def filter_data(rawData):
# Define columns to drop
columns_to_drop = ['Customer Type', 'Country']
# Drop the specified columns in PySpark DataFrame
filteredRetailData = rawData.drop(*columns_to_drop)
# Log the data after dropping columns
unique_file_name = 'script://Weekly_Sales/Filtered_df'
with dataset_op_logger(unique_file_name, "read", with_schema=True, with_preview=True) as logger:
logger.set(data=filteredRetailData)
return filteredRetailData
@task
def write_data_by_product_line(filteredData):
# Filter data for Camping Equipment and write to CSV
campingEquipment = filteredData.filter(col('Product line') == 'Camping Equipment')
campingEquipment.write.csv("Camping_Equipment.csv", header=True, mode="overwrite")
# Log writing the Camping Equipment CSV
log_dataframe("camping_equipment", campingEquipment, with_schema=True, with_stats=True)
# Filter data for Golf Equipment and write to CSV
golfEquipment = filteredData.filter(col('Product line') == 'Golf Equipment')
golfEquipment.write.csv("Golf_Equipment.csv", header=True, mode="overwrite")
# Log writing the Golf Equipment CSV
log_dataframe("golf_equipment", golfEquipment, with_schema=True, with_stats=True)
def prepare_retail_data():
with dbnd_tracking(
conf={
"tracking": {
"track_source_code": True
},
"log": {
"preview_head_bytes": 15360,
"preview_tail_bytes": 15360
}
}
):
logger.info("Running Databand spark application!")
start_time_milliseconds = int(round(time.time() * 1000))
log_metric("metric_check", "OK")
# Call the step job - create sample data
rawData = create_sample_data()
# Filter data
filteredData = filter_data(rawData)
# Write data by product line
write_data_by_product_line(filteredData)
end_time_milliseconds = int(round(time.time() * 1000))
elapsed_time = end_time_milliseconds - start_time_milliseconds
log_metric('elapsed-time', elapsed_time)
logger.info(f"Total pipeline running time: {elapsed_time:.2f} milliseconds")
logger.info("Spark execution completed..")
log_metric("is-success", "OK")
# Invoke the main function
prepare_retail_data()
Breve panoramica delle API Databand e delle funzionalità utilizzate:
-
dbnd_tracking: Inizializza il tracciamento della pipeline o dell'applicazione, configurando le impostazioni di Databand e registrando i dettagli dell'esecuzione. Uso:with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"): # Pipeline code -
task: Contrassegna una funzione come attività Databand, consentendo di tracciare e monitorare le singole fasi della pipeline.Uso:
@task def my_task_function(): -
dataset_op_logger: Registra le operazioni sugli insiemi di dati, compresi lo schema e le statistiche.Uso:
with dataset_op_logger(dataset_name, operation_type) as logger: logger.set(data=my_dataframe) -
log_metric: Registra metriche personalizzate per monitorare le prestazioni o altri dati quantitativi durante l'esecuzione.Uso:
log_metric("metric_name", metric_value) -
log_dataframe: Registra i dettagli di DataFrame,, come lo schema e le statistiche, per monitorare le trasformazioni dei dati.Uso:
log_dataframe("dataframe_name", my_dataframe, with_schema=True, with_stats=True)
Per informazioni sull'invio di lavori Spark, vedere Invio di applicazioni del motore.
Dopo aver inviato la domanda Spark, riceverete un messaggio di conferma con l'ID della domanda e la versione di Spark. Conservate queste informazioni per monitorare lo stato di esecuzione del lavoro Spark inviato. È possibile monitorare e seguire gli insiemi di dati utilizzando le funzioni di monitoraggio di Databand all'interno dell'ambiente Databand.
Per visualizzare il cruscotto Databand per il monitoraggio, andare su Configurazioni > IBM Osservabilità dei dati per Databand e fare clic su Visualizza Databand.
Per ulteriori informazioni ed esempi, vedere:
- IBM Osservabilità dei dati per banda dati
- Per le applicazioni di PySpark: Tracciamento PySpark
- Per Spark( Java / Scala ): Tracciamento di Spark(Scala / Java)