Suivi de l'exécution des applications Spark à l'aide de Databand
S'applique à: Moteur à étincelles
L'intégration de Databand avec Spark améliore les capacités de surveillance en fournissant des informations qui vont au-delà de l' Spark UI et de l'historique Spark.
Databand améliore la surveillance des applications Spark grâce aux éléments suivants :
- Surveillance avancée: Les annotations de tâches de Databand vous permettent de marquer et de suivre les étapes cruciales de votre application Spark, offrant un niveau de surveillance plus significatif par rapport aux jobs, étapes ou tâches Spark.
- Suivi des ensembles de données: Databand surveille les ensembles de données qui sont accédés et modifiés pendant l'exécution de votre application Spark, offrant ainsi une meilleure visibilité sur vos flux de données.
- Alertes personnalisées: vous pouvez configurer des alertes pour des étapes spécifiques de votre application ou suivre des métriques clés de l'ensemble des données, ce qui vous permet d'identifier et de traiter les problèmes potentiels à un stade précoce.
Pour commencer à utiliser Databand, vous devez avoir un abonnement Databand actif. Vous pouvez l'obtenir en demandant une instance de l'application en nuage Databand (SaaS), qui est déployée par l'équipe Databand, ou en optant pour une installation auto-hébergée (sur site). Pour intégrer la bande de données à votre instance watsonx.data, vous devez disposer des informations d'identification suivantes :
- Adresse de l'environnement: L' URL de votre environnement Databand (exemple : yourcompanyname.databand.ai ).
- Code d'accès: Un jeton d'accès Databand nécessaire pour se connecter à l'environnement. Vous pouvez générer et gérer des jetons via l'interface utilisateur de Databand si nécessaire. Pour plus de détails, consultez le site : Gestion des jetons d'accès personnels.
Procédure
- Connectez vous à la console watsonx.data.
- Dans le menu de navigation, allez à Configurations > IBM Data Observability by Databand.
- Saisissez l'adresse de l'environnement dans le champ URL et le jeton d'accès dans le champ Jeton d'accès.
- Cliquez sur Tester la connexion pour valider la connexion.
- Cliquez sur Enregistrer pour sauvegarder les détails. Vous pouvez modifier les détails en cliquant sur Modifier.
Databand prend effet pour les nouveaux jobs Spark qui sont lancés après l'avoir activé. Les emplois précédents et en cours ne sont pas enregistrés à des fins d'analyse.
Auditeur de bande de données
Cette méthode permet de suivre automatiquement les opérations effectuées sur les ensembles de données. Votre script Spark peut bénéficier d'un suivi automatique des opérations sur les ensembles de données.
Décorateurs de bande de données et API de journalisation
Pour utiliser cette méthode, vous devez importer le module dbnd, ce qui nécessite des modifications du code.
Exemple de bout en bout d'utilisation des API Databand
L'exemple suivant illustre l'utilisation des API 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()
Brève présentation des API de la bande de données et des fonctionnalités utilisées :
-
dbnd_tracking: Initialise le suivi de votre pipeline ou de votre application, en configurant les paramètres Databand et en enregistrant les détails de l'exécution. Syntaxe :with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"): # Pipeline code -
task: Marque une fonction en tant que tâche Databand, ce qui permet de suivre et de contrôler les différentes étapes de votre pipeline.Syntaxe :
@task def my_task_function(): -
dataset_op_logger: Enregistre les opérations sur les ensembles de données, y compris le schéma et les statistiques.Syntaxe :
with dataset_op_logger(dataset_name, operation_type) as logger: logger.set(data=my_dataframe) -
log_metric: Enregistre des mesures personnalisées pour suivre les performances ou d'autres données quantitatives pendant l'exécution.Syntaxe :
log_metric("metric_name", metric_value) -
log_dataframe: Enregistre les détails d'un site DataFrame,, tels que le schéma et les statistiques, afin de surveiller les transformations de données.Syntaxe :
log_dataframe("dataframe_name", my_dataframe, with_schema=True, with_stats=True)
Pour plus d'informations sur la soumission des tâches Spark, voir Soumettre des demandes de moteur.
Après avoir soumis la demande Spark, vous recevrez un message de confirmation avec l'identifiant de la demande et la version de Spark. Conservez ces informations pour suivre l'état d'exécution de votre job Spark soumis. Vous pouvez surveiller et suivre les ensembles de données en utilisant les fonctions de suivi de Databand dans l'environnement Databand.
Pour afficher le tableau de bord de la bande de données pour le suivi, allez dans Configurations > IBM Data Observability by Databand (Observabilité des données par bande de données ) et cliquez sur View Databand (Afficher la bande de données).
Pour plus d'informations et d'exemples, voir :
- IBM Observabilité des données par bande de données
- Pour PySpark Applications : Suivi PySpark
- Pour Spark( Java / Scala ): Suivi de l'étincelle(Scala / Java)