Monitorización de la ejecución de aplicaciones Spark mediante Databand
Se aplica a: Motor de chispa
La integración de Databand con Spark mejora las capacidades de monitorización al proporcionar información que va más allá de Spark UI y Spark History.
Databand mejora la monitorización de aplicaciones Spark de la siguiente manera:
- Monitorización avanzada: Las anotaciones de tareas de Databand le permiten etiquetar y rastrear etapas cruciales de su aplicación Spark, ofreciendo un nivel de monitorización más significativo en comparación con los trabajos, etapas o tareas de Spark.
- Seguimiento de conjuntos de datos: Databand monitoriza los conjuntos de datos a los que se accede y que se modifican durante la ejecución de su aplicación Spark, proporcionando una visibilidad mejorada de sus flujos de datos.
- Alertas personalizadas: puede configurar alertas para etapas específicas de su aplicación o realizar un seguimiento de las métricas clave del conjunto de datos, lo que le permite identificar y abordar posibles problemas con antelación.
Para empezar a utilizar Databand, debe tener una suscripción activa a Databand. Puede obtenerla solicitando una instancia de la aplicación en la nube de Databand (SaaS), que despliega el equipo de Databand, u optando por una instalación autoalojada (on-premises). Para integrar databand con su instancia watsonx.data, debe disponer de las siguientes credenciales:
- Dirección del entorno: La URL de tu entorno Databand (ejemplo: yourcompanyname.databand.ai ).
- Token de acceso: Un token de acceso a Databand que se necesita para conectarse al entorno. Puede generar y gestionar tokens a través de la interfaz de usuario de Databand según sea necesario. Para más detalles, visite: Gestión de tokens de acceso personales.
Procedimiento
- Inicie la sesión en la consola de watsonx.data.
- En el menú de navegación, vaya a Configuraciones > IBM Data Observability by Databand.
- Introduzca la dirección del entorno en el campo URL y el código de acceso en el campo Código de acceso.
- Haga clic en Probar conexión para validar la conexión.
- Haga clic en Guardar para guardar los detalles. Puede editar los detalles haciendo clic en Editar.
Databand tiene efecto para los nuevos trabajos de Spark que se inicien después de activarlo. Los trabajos anteriores y en curso no se registran para su análisis.
Escucha de banda de datos
Este método rastrea automáticamente las operaciones del conjunto de datos. Su script Spark puede beneficiarse del seguimiento automático de las operaciones del conjunto de datos.
Decoradores de bases de datos y API de registro
Para utilizar este método, debe importar el módulo dbnd, lo que requiere modificaciones en el código.
Ejemplo de uso integral de las API de Databand
El siguiente ejemplo demuestra el uso de las API de 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 descripción de las API de Databand y de las funciones utilizadas:
-
dbnd_tracking: Inicializa el seguimiento de su pipeline o aplicación, configurando los ajustes de Databand y registrando los detalles de ejecución. Uso:with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"): # Pipeline code -
task: Marca una función como tarea de Databand, lo que permite el seguimiento y la supervisión de pasos individuales en su pipeline.Uso:
@task def my_task_function(): -
dataset_op_logger: Registra las operaciones realizadas en los conjuntos de datos, incluidos el esquema y las estadísticas.Uso:
with dataset_op_logger(dataset_name, operation_type) as logger: logger.set(data=my_dataframe) -
log_metric: Registra métricas personalizadas para realizar un seguimiento del rendimiento u otros datos cuantitativos durante la ejecución.Uso:
log_metric("metric_name", metric_value) -
log_dataframe: Registra detalles sobre DataFrame,, como el esquema y las estadísticas, para supervisar las transformaciones de datos.Uso:
log_dataframe("dataframe_name", my_dataframe, with_schema=True, with_stats=True)
Para obtener información sobre el envío de trabajos de Spark, consulte Enviar solicitudes de motores.
Tras enviar la solicitud de Spark, recibirá un mensaje de confirmación con el ID de la solicitud y la versión de Spark. Conserve esta información para realizar un seguimiento del estado de ejecución de su trabajo Spark enviado. Puede supervisar y realizar un seguimiento de los conjuntos de datos utilizando las funciones de seguimiento de Databand dentro del entorno de Databand.
Para ver el cuadro de mando de Databand para el seguimiento, vaya a Configuraciones > IBM Observabilidad de datos por Databand y haga clic en Ver Databand.
Para más información y ejemplos, consulte:
- IBM Observabilidad de los datos por banda de datos
- Para aplicaciones PySpark: Seguimiento PySpark
- Para Spark( Java / Scala ): Seguimiento de Spark(Scala / Java)