Databandを使用してSparkアプリケーションの実行を監視する

適用対象スパークエンジン

DatabandとSparkの統合は、 Spark UI Spark Historyを超えた洞察を提供することで、モニタリング機能を強化します。

DatabandはSparkアプリケーションのモニタリングを以下のように改善する:

  • 高度なモニタリング:Databand のタスクアノテーションは、Spark アプリケーションの重要なステージにタグを付けて追跡することを可能にし、Spark のジョブ、ステージ、またはタスクと比較して、より有意義なレベルのモニタリングを提供します。
  • データセットのトラッキング:Databand は、Spark アプリケーションの実行中にアクセスおよび変更されたデータセットを監視し、データフローの可視性を高めます。
  • カスタム アラート: アプリケーションの特定の段階にアラートを設定したり、主要なデータ セット メトリックを追跡したりすることができます。

Databandを始めるには、アクティブなDatabandサブスクリプションが必要です。 Databand チームによってデプロイされる Databand クラウドアプリケーション (SaaS) インスタンスをリクエストするか、またはセルフホスト (オンプレミス) インストールを選択することで、これを取得できます。 databandと watsonx.data インスタンスを統合するには、以下の認証情報が必要です:

  • 環境アドレス :Databand環境の URL (例: yourcompanyname.databand.ai )。
  • アクセストークン:環境に接続するために必要な Databand アクセストークン。 必要に応じて Databand UI からトークンを生成し、管理することができます。 詳細はこちらをご覧ください:個人アクセストークンの管理 をご覧ください。

手順

  1. watsonx.data コンソールにログインします。
  2. ナビゲーション・メニューから、 Configurations > IBM Data Observability by Databand に進みます。
  3. URL フィールドに環境アドレス、 アクセストークンフィールドにアクセストークンを入力します。
  4. Test connectionをクリックして接続を確認する。
  5. Saveをクリックして詳細を保存します。 Edit(編集) 」をクリックすると、詳細を編集できます。

Databandは、有効化後に開始される新しいSparkジョブに対して有効になる。 前職や進行中の仕事は分析のために記録されない。

データバンドリスナー

この方法は、データセットの操作を自動的に追跡する。 Sparkスクリプトは、データセット操作の自動追跡の恩恵を受けることができます。

データバンド・デコレーターとロギングAPI

この方法を使うには、 dbnd モジュールをインポートする必要があり、コードの修正が必要になる。

Databand 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()

Databand APIの概要と使用機能:

  • dbnd_tracking:パイプラインやアプリケーションのトラッキングを初期化し、Databand の設定を行い、実行の詳細をログに記録します。 使用方法:

    with dbnd_tracking(conf={...}, job_name="job_name", run_name="run_name"):
     # Pipeline code
    
  • task:機能を Databand タスクとしてマークし、パイプラインの個々のステップの追跡と監視を可能にします。

    使用方法:

    @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)
    

Sparkジョブの送信については、 エンジンアプリケーションの送信を 参照してください。

Sparkアプリケーションを送信すると、アプリケーションIDとSparkバージョンが記載された確認メッセージが届きます。 この情報は、投入したSparkジョブの実行状況を追跡するために保管してください。 Databand環境内でDatabandのトラッキング機能を使用することで、データセットを監視・追跡することができます。

追跡用のデータバンド・ダッシュボードを表示するには、 Configurations > IBM Data Observability by Databandに進み、 View Databandをクリックする。

詳細と例については、こちらを参照: