Überwachung der Spark-Anwendungsläufe mit Databand
Gilt für: Funkenmotor
Die Databand-Integration mit Spark verbessert die Überwachungsfunktionen, indem sie Einblicke bietet, die über Spark UI und Spark History hinausgehen.
Databand verbessert die Überwachung von Spark-Anwendungen wie folgt:
- Erweiterte Überwachung: Die Task-Annotationen von Databand ermöglichen es Ihnen, wichtige Phasen Ihrer Spark-Anwendung zu markieren und zu verfolgen, was im Vergleich zu Spark-Jobs, Phasen oder Tasks ein aussagekräftigeres Niveau der Überwachung bietet.
- Datensatzverfolgung: Databand überwacht die Datensätze, auf die während des Laufs Ihrer Spark-Anwendung zugegriffen wird und die geändert werden, und bietet so einen besseren Einblick in Ihre Datenflüsse.
- Benutzerdefinierte Warnmeldungen: Sie können Warnmeldungen für bestimmte Phasen Ihrer Anwendung konfigurieren oder wichtige Datensatzmetriken verfolgen, so dass Sie potenzielle Probleme frühzeitig erkennen und beheben können.
Um mit Databand arbeiten zu können, müssen Sie über ein aktives Databand-Abonnement verfügen. Sie können entweder eine Databand Cloud-Anwendung (SaaS) anfordern, die vom Databand-Team bereitgestellt wird, oder sich für eine selbst gehostete (Vor-Ort-)Installation entscheiden. Für die Integration von databand in Ihre watsonx.data benötigen Sie die folgenden Anmeldedaten:
- Adresse der Umgebung: Die URL für Ihre Databand-Umgebung (Beispiel: yourcompanyname.databand.ai ).
- Zugangs-Token: Ein Databand-Zugangstoken, der für die Verbindung mit der Umgebung benötigt wird. Sie können Token bei Bedarf über die Databand-Benutzeroberfläche erstellen und verwalten. Weitere Einzelheiten finden Sie unter: Verwaltung von persönlichen Zugangstokens.
Vorgehensweise
- Melden Sie sich bei der watsonx.data-Konsole an.
- Gehen Sie im Navigationsmenü zu Konfigurationen > IBM Data Observability by Databand.
- Geben Sie die Umgebungsadresse in das Feld URL und das Zugriffstoken in das Feld Zugriffstoken ein.
- Klicken Sie auf Verbindung testen, um die Verbindung zu überprüfen.
- Klicken Sie auf Speichern, um die Details zu speichern. Sie können die Details bearbeiten, indem Sie auf Bearbeiten klicken.
Databand wird für neue Spark-Aufträge wirksam, die nach seiner Aktivierung gestartet werden. Frühere und laufende Aufträge werden für die Analyse nicht erfasst.
Databand-Hörer
Mit dieser Methode werden Datensatzoperationen automatisch verfolgt. Ihr Spark-Skript kann von der automatischen Verfolgung von Datensatzoperationen profitieren.
Databand-Dekoratoren und Protokollierungs-API
Um diese Methode zu verwenden, müssen Sie das Modul dbnd importieren, was Codeänderungen erfordert.
End-to-End-Beispiel für die Verwendung von Databand-APIs
Das folgende Beispiel zeigt die Verwendung von dbnd APIs.
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()
Kurzer Überblick über die verwendeten Databand APIs und Funktionen:
-
dbnd_tracking: Initialisiert die Verfolgung für Ihre Pipeline oder Anwendung, konfiguriert die Databand-Einstellungen und protokolliert Ausführungsdetails. Verwendung:with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"): # Pipeline code -
task: Markiert eine Funktion als Databand-Aufgabe und ermöglicht so die Verfolgung und Überwachung einzelner Schritte in Ihrer Pipeline.Verwendung:
@task def my_task_function(): -
dataset_op_logger: Protokolliert Operationen auf Datensätzen, einschließlich Schema und Statistiken.Verwendung:
with dataset_op_logger(dataset_name, operation_type) as logger: logger.set(data=my_dataframe) -
log_metric: Zeichnet benutzerdefinierte Metriken auf, um die Leistung oder andere quantitative Daten während der Ausführung zu verfolgen.Verwendung:
log_metric("metric_name", metric_value) -
log_dataframe: Protokolliert Details über eine DataFrame, wie Schema und Statistiken zur Überwachung von Datentransformationen.Verwendung:
log_dataframe("dataframe_name", my_dataframe, with_schema=True, with_stats=True)
Informationen zum Einreichen von Spark-Aufträgen finden Sie unter Einreichen von Engine-Anwendungen.
Nachdem Sie den Spark-Antrag eingereicht haben, erhalten Sie eine Bestätigungsnachricht mit der Antrags-ID und der Spark-Version. Bewahren Sie diese Informationen auf, um den Ausführungsstatus Ihres eingereichten Spark-Auftrags zu verfolgen. Sie können Datensätze mit Hilfe der Tracking-Funktionen von Databand innerhalb der Databand-Umgebung überwachen und verfolgen.
Um das Databand-Dashboard für die Nachverfolgung anzuzeigen, gehen Sie zu Konfigurationen > IBM Data Observability by Databand und klicken Sie auf View Databand.
Weitere Informationen und Beispiele finden Sie unter:
- IBM Beobachtbarkeit der Daten nach Datenband
- Für PySpark Anwendungen: Verfolgung PySpark
- Für Spark( Java / Scala ): Verfolgung von Spark(Scala / Java)