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

  1. Accedi alla console watsonx.data.
  2. Dal menu di navigazione, andare su Configurazioni > IBM Data Observability by Databand.
  3. Inserire l'indirizzo dell'ambiente nel campo URL e il token di accesso nel campo Token di accesso.
  4. Fare clic su Prova connessione per convalidare la connessione.
  5. 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: