데이터밴드를 사용하여 Spark 애플리케이션 실행 모니터링
에 적용됩니다: 스파크 엔진
데이터밴드와 Spark의 통합은 Spark UI 및 Spark 히스토리를 넘어서는 인사이트를 제공함으로써 모니터링 기능을 향상시킵니다.
데이터밴드는 다음과 같은 방식으로 Spark 애플리케이션 모니터링을 개선합니다:
- 고급 모니터링: 데이터밴드의 작업 주석을 사용하면 Spark 애플리케이션의 중요한 단계에 태그를 지정하고 추적할 수 있어 Spark 작업, 단계 또는 작업에 비해 더 의미 있는 수준의 모니터링을 제공합니다.
- 데이터 세트 추적: 데이터밴드는 Spark 애플리케이션 실행 중에 액세스 및 수정되는 데이터세트를 모니터링하여 데이터 흐름에 대한 향상된 가시성을 제공합니다.
- 사용자 지정 알림: 애플리케이션의 특정 단계에 대한 알림을 구성하거나 주요 데이터 세트 메트릭을 추적하여 잠재적인 문제를 조기에 파악하고 해결할 수 있습니다.
데이터밴드를 시작하려면 활성 상태의 데이터밴드 구독이 있어야 합니다. 데이터밴드 팀에서 배포하는 데이터밴드 클라우드 애플리케이션(SaaS) 인스턴스를 요청하거나 자체 호스팅(온프레미스) 설치를 선택하면 이 인스턴스를 얻을 수 있습니다. 데이터밴드를 watsonx.data 인스턴스와 통합하려면 다음 자격 증명이 있어야 합니다:
- 환경 주소: 데이터밴드 환경의 URL (예: yourcompanyname.databand.ai )입니다.
- 액세스 토큰: 환경에 연결하는 데 필요한 데이터밴드 액세스 토큰입니다. 필요에 따라 데이터밴드 UI를 통해 토큰을 생성하고 관리할 수 있습니다. 자세한 내용은 여기를 참조하세요: 개인 액세스 토큰 관리하기.
프로시저
- watsonx.data 콘솔에 로그인하십시오.
- 탐색 메뉴에서 구성 > 데이터밴드별 IBM 데이터 관측성으로 이동합니다.
- URL 필드에 환경 주소를 입력하고 액세스 토큰 필드에 액세스 토큰을 입력합니다.
- 연결 테스트를 클릭하여 연결의 유효성을 검사합니다.
- 저장을 클릭하여 세부 정보를 저장합니다. 편집을 클릭하여 세부 정보를 편집할 수 있습니다.
데이터밴드는 활성화한 후 시작하는 새 Spark 작업에 적용됩니다. 이전 작업과 진행 중인 작업은 분석을 위해 기록되지 않습니다.
데이터밴드 리스너
이 방법은 데이터 세트 작업을 자동으로 추적합니다. Spark 스크립트는 데이터 세트 작업의 자동 추적을 통해 이점을 얻을 수 있습니다.
데이터밴드 데코레이터 및 로깅 API
이 방법을 사용하려면 코드를 수정해야 하는 dbnd 모듈을 가져와야 합니다.
데이터밴드 API 사용의 엔드투엔드 예시
다음 예는 dbnd API의 사용법을 보여줍니다.
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()
데이터밴드 API 및 사용 기능에 대한 간략한 개요:
-
dbnd_tracking: 파이프라인 또는 애플리케이션에 대한 추적을 초기화하여 데이터밴드 설정을 구성하고 실행 세부 정보를 기록합니다. 사용법:with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"): # Pipeline code -
task: 함수를 데이터밴드 작업으로 표시하여 파이프라인의 개별 단계를 추적하고 모니터링할 수 있도록 합니다.사용법:
@task def my_task_function(): -
dataset_op_logger: 스키마 및 통계를 포함하여 데이터 집합에 대한 작업을 기록합니다.사용법:
with dataset_op_logger(dataset_name, operation_type) as logger: logger.set(data=my_dataframe) -
log_metric: 실행 중 성능 또는 기타 정량적 데이터를 추적하기 위해 사용자 지정 지표를 기록합니다.사용법:
log_metric("metric_name", metric_value) -
log_dataframe: 데이터 변환을 모니터링하기 위해 스키마 및 통계와 같은 DataFrame, 에 대한 세부 정보를 기록합니다.사용법:
log_dataframe("dataframe_name", my_dataframe, with_schema=True, with_stats=True)
스파크 작업 제출에 대한 자세한 내용은 엔진 신청서 제출을 참조하세요.
스파크 신청서를 제출하면 신청서 ID와 스파크 버전이 포함된 확인 메시지를 받게 됩니다. 제출한 스파크 작업의 실행 상태를 추적하기 위해 이 정보를 보관하세요. 데이터밴드 환경 내에서 데이터밴드의 추적 기능을 사용하여 데이터세트를 모니터링하고 추적할 수 있습니다.
추적용 데이터밴드 대시보드를 보려면 구성 > IBM 데이터밴드별 데이터 관찰 가능성으로 이동하여 데이터밴드 보기를 클릭합니다.
자세한 정보와 예시는 다음을 참조하세요:
- IBM 데이터 밴드별 데이터 가시성
- PySpark 애플리케이션용: 추적 PySpark
- 스파크( Java / Scala ): 스파크 추적(Scala / Java)