Azkabanプラグイン開発の実践:カスタムJobタイプとViewer拡張の実装手法

Azkabanプラグインアーキテクチャの概要

Azkabanのプラグインシステムは、モジュール化と高い拡張性を実現するように設計されています。主に「Jobタイプ」と「Viewer」の2つのカテゴリに分類されます。Jobタイププラグインはタスク実行機能を拡張し、ViewerプラグインはHDFS上のファイル閲覧機能を強化します。これらはJavaベースで実装され、特定のインターフェースとプロパティファイルを通じてシステムに統合されます。コア機能の安定性を損なうことなく、独立して開発・デプロイできるのが特徴です。

開発環境とプロジェクト構成

プラグイン開発を開始する前に、以下の環境を準備する必要があります。

  • Java JDK 8以上
  • MavenまたはGradleビルドツール
  • Azkabanソースコード

標準的なプラグインプロジェクトのディレクトリ構成は以下のようになります。


azkaban-extension/
├── src/
│   ├── main/
│   │   ├── java/
│   │   │   └── com/
│   │   │       └── example/
│   │   │           ├── jobtype/    # Jobタイプ拡張コード
│   │   │           └── viewer/     # Viewer拡張コード
│   │   └── resources/
│   │       └── plugin.properties   # プラグイン定義ファイル
└── build.gradle                    # ビルド設定ファイル

カスタムJobタイププラグインの実装

基本設計と実装ステップ

Jobタイププラグインは、新しいタスク実行ロジックを定義するために使用されます。既存のJavaProcessJobProcessJobを継承して実装します。

まず、Javaプロセスを実行するJobクラスを作成します。


public class DataProcessingJob extends JavaProcessJob {
    public DataProcessingJob(String jobId, Props systemProps, Props jobProps, Logger logger) {
        super(jobId, systemProps, new Props(systemProps, jobProps), logger);
    }
    
    @Override
    protected String getJavaClass() {
        return "org.example.data.BatchProcessor";
    }
    
    // 必要に応じて他のメソッドをオーバーライド
}

次に、src/main/resourcesディレクトリにplugin.propertiesを作成し、クラスパスとJobタイプ名をマッピングします。


job.class=com.example.jobtype.DataProcessingJob
jobtype.name=data_processor

实例:ProcessBuilderを利用したPython実行Job

外部スクリプトを実行するJobタイプを実装する場合は、ProcessJobを継承し、ProcessBuilderを使用してプロセスを制御すると柔軟性が高まります。


public class PythonExecutionJob extends ProcessJob {
    public PythonExecutionJob(String id, Props sysProps, Props jobProps, Logger log) {
        super(id, sysProps, jobProps, log);
    }
    
    @Override
    public void run() throws Exception {
        String targetScript = getJobProps().getString("target.script");
        String pythonEnv = getJobProps().getString("python.env", "python3");
        
        ProcessBuilder pb = new ProcessBuilder(pythonEnv, targetScript);
        pb.redirectErrorStream(true);
        pb.environment().putAll(getJobProps().getFlattened());
        
        Process process = pb.start();
        int exitStatus = process.waitFor();
        
        if (exitStatus != 0) {
            throw new RuntimeException("Python script execution failed. Exit code: " + exitStatus);
        }
    }
}

Viewerプラグインの実装

HDFSファイルビューアの拡張

Viewerプラグインは、AzkabanのWeb UI上で特定のファイル形式をレンダリングするために使用されます。HdfsFileViewerを継承し、サポートする拡張子と読み取りロジックを定義します。

以下は、JSONファイルをフォーマットして表示するViewerの実装例です。


public class JsonFileViewer extends HdfsFileViewer {
    private static final String DISPLAY_NAME = "JSON Viewer";
    private final Set<String> supportedExtensions = new HashSet<>(Arrays.asList(".json", ".jsonl"));

    @Override
    public String getName() {
        return DISPLAY_NAME;
    }
    
    @Override
    public Set<Capability> getCapabilities(FileSystem fs, Path path) {
        String extension = path.getName().substring(path.getName().lastIndexOf('.'));
        if (supportedExtensions.contains(extension.toLowerCase())) {
            return EnumSet.of(Capability.READ);
        }
        return EnumSet.noneOf(Capability.class);
    }
    
    @Override
    public void displayFile(FileSystem fs, Path path, OutputStream out, int start, int end) throws IOException {
        // JSONストリームの読み込み、パース、およびHTMLレンダリング処理をここに実装
    }
}

対応するプロパティファイルは以下のようになります。


viewer.class=com.example.viewer.JsonFileViewer
viewer.name=json_renderer

設定とデプロイメント

プラグインをAzkabanサーバーに統合するには、ビルドしたアーティファクトを適切なディレクトリに配置する必要があります。

  1. プロジェクトをビルドし、依存関係を含むFat JAR(または配布用ZIP)を生成します。
  2. Azkabanインストールディレクトリ配下のplugins/jobtypes/[jobtype_name]またはplugins/viewers/[viewer_name]ディレクトリを作成します。
  3. 生成したJARファイルとplugin.propertiesを対象ディレクトリにコピーします。
  4. クラスパスに外部ライブラリが必要な場合は、lib/サブディレクトリに配置し、プロパティファイルでclasspath=lib/*と指定します。
  5. Azkaban ExecutorおよびWebサーバーを再起動して変更を反映させます。

テストとデバッグ

単体テストの設計

プラグインのロジックが正しく動作することを保証するため、モックオブジェクトを活用した単体テストを記述します。


public class DataProcessingJobTest {
    @Test
    public void verifyJobExecution() throws Exception {
        Props sysProps = new Props();
        Props jobProps = new Props();
        jobProps.put("job.class", "org.example.data.BatchProcessor");
        jobProps.put("target.script", "dummy.py");
        
        Logger mockLogger = LogManager.getLogger(DataProcessingJobTest.class);
        DataProcessingJob job = new DataProcessingJob("test-job-01", sysProps, jobProps, mockLogger);
        
        // 実行結果の検証とアサーション
        job.run();
    }
}

デバッグ手法

  • ログ出力の強化: プラグイン内の重要な分岐点でLoggerを用いて詳細なDEBUGレベルのログを出力します。
  • リモートデバッグ: Azkaban ExecutorのJVM起動オプションに-agentlib:jdwp=transport=dt_socket,server=y,suspend=n,address=5005を追加し、IDEからリモートデバッガーをアタッチします。
  • 実行ログの分析: azkaban-execserver/logsディレクトリ内のazkaban-execserver.logと、各Jobの実行ディレクトリに生成される.logファイルを確認します。

高度な拡張とトラブルシューティング

イベントトリガーとUIカスタマイズ

Azkabanはフローの依存関係に基づいたイベントトリガーメカニズムをサポートしています。azkaban.flowtriggerパッケージ内のインターフェースを実装することで、外部システムからのWebhookやデータベースの状態変更をトリガーにしたJob実行が可能になります。また、Web UIのカスタマイズが必要な場合は、azkaban-web-serverモジュールのフロントエンドリソース(Less/JavaScript)を拡張し、カスタムViewer用の専用のレンダリングコンポーネントを追加できます。

一般的な問題と解決策

  • プラグインがロードされない: plugin.propertiesjob.classまたはviewer.classのFQCN(完全修飾クラス名)にタイポがないか確認してください。また、JARファイルに必要な依存クラスがすべて含まれているか(Fat JAR化されているか)を検証します。
  • クラスキャスト例外 (ClassCastException): Azkabanコアライブラリとプラグイン側で異なるバージョンのライブラリ(例:GuavaやJackson)が競合している可能性があります。Maven/Gradleのprovidedスコープを利用して、コア側とバージョンを合わせるか、 shading(リロケーション)を適用してください。
  • HDFS権限エラー: ViewerプラグインがHDFS上のファイルにアクセスする際、Azkaban Executorを実行しているOSユーザーに適切なHDFS読み取り権限が付与されているか確認してください。Kerberos環境では、チケットの更新ロジックが正しく機能しているかも検証対象となります。

タグ: Azkaban Java Hadoop HDFS WorkflowManager

9月11日 02:02 投稿