データのパイプラインにおける自動テストの重要な役割

Apache Spark のパワー ミッション クリティカルな分析、機械学習ワークフロー、リアルタイムの意思決定に基づいて構築されたデータパイプライン。 トランスフォーメーションにおける単一のロジックエラーでさえ、ダウンストリームレポートを強制的に中断したり、業務の不適切な操作をトリガーしたり、高価な計算リソースを無駄にしたりすることができます。 マニュアルテスト - いくつかの行をチェックアウトしたり、データの一部をサブセットにスクリプトを実行したりするなど、現代のエンジニアリングデータパイプラインの複雑さと速度を低下させません。 自動化されたテストフレームワークは、システムの状態によってこのギャップに対処し、その結果を検証し、正確な作業を削減し、その結果を検証します。

Spark パイプラインのテストフレームワークの設計

Spark の強力なテストフレームワークは、データパイプライン開発のアートを繰り返し可能なエンジニアリングの分野に変換します。フレームワークは、モジュール式で再利用可能なコンポーネントに、ユニット、統合、エンドツーエンドのテスト用に構成できるコンポーネントを分離する必要があります。以下は、重要なビルディングブロックです。

データの生成をテストする

代表的なテストデータは、効果的なテストの基礎です。大きめの、しばしば敏感で、保守が困難である、生産表全体をコピーする代わりに、境界条件、null値、重複キー、および予期しないフォーマットを練習する、小さな集中したデータセットを作成します。Sparkの組み込み[]を使用して、決定的な入力を生成します。より複雑なシナリオでは、ランダムなが生成する工場やビルダーをレバレッジしますが、ランダムな合成データを生成する[FLT]を[FLT]]]を[FLT]]]]を[FLT]]]]に置き換えてください。

試験ケースと認証

各テストケースは、特定の入力状態を定義し、変換または一連の変換を実行し、出力に対するアサーションを適用します。 一般的なアサーションパターンは次のとおりです。

  • [] 降下レベル平等:[ 想定されたデータフレームの列ごとに比較します。
  • [] シュマバリデーション:[]]]] 出力スキーマが意図したタイプと無効なプロパティにマッチすることを確認します。
  • []集合チェック:[グループごとの操作後のカウント、合計、またはユニークな値を確認します。
  • [] 業務執行:]] は、派生した列(例えば、年齢のバケツ、異常フラグ)が許容範囲内で落ちることを確認します。

明確で自己文書化ステートメントとしてアサーションを書きます。 ScalaTest では または を使用します。 PyTest では、パンダの互換性のあるアサーションまたは専用の ]] と組み合わせます。chisui/assert-spark ライブラリ。

実行環境

ローカルモードでテストをスパークしてクラスターのオーバーヘッドを回避します。 を ]] で設定します。 複数の実行を 1 つの JVM または Python プロセスで実行します。 並列を低数に設定します(例: )) テスト時間を短縮します。 Scala プロジェクトの場合、]] から Spark test ライブラリベース ライブラリのレイキャリブ を 1 つのセッションごとに、 S] をクリーンにテストします。

検証とレポート

自動テスト実行では、ログ、パス/失敗数、エラーの詳細が生成されます。テストレポートを継続的に統合(CI)ダッシュボードに統合することで、チームメンバーはパイプラインコンポーネントがどのコンポーネントが壊れたのかを素早く特定できます。 のようなツール]またはScalaTestとPyTestの組み込みXMLレポーターは、入力データを表示したり、実際の結果や実行期間を想定したバロース可能なレポートを生成したりできます。 この透明性は、root-cause分析を加速し、品質を向上します。

実用的な実装戦略

フレームワークコンポーネントを現実世界スパークパイプラインのテストシナリオにマップする次のアプローチ。

ユニットテストの変換

ユニットテストは、DataFrame を操作する単一の関数またはメソッドを検証します。例えば、タイムスタンプ文字列をクリーンにする関数を考慮する: 。ユニットテストは、有効で、malformed、null タイムスタンプを持つ小さな DataFrame を作成し、関数を呼び出し、出力カラムがその列の期待値だけを含んでいることを主張します。テストはローカルモードで実行し、数行だけを処理するため、すべての開発者がテストを行なうために、すべての開発者がテストを行に渡します。

統合テスト

統合テストでは、複数の変換が正しく機能していることが確認されます。例えば、パイプラインは、生のJSONイベント、フラットなネスト構造を読み込み、寸法表に結合し、ウィンドウ機能を適用することができます。統合テストは、すべてのソースデータ(または現実的な合成代替物)をロードし、ジョブのロジック全体を特定のステージまで実行し、そのステージの出力が既知のゴールデンデータセットにマッチすることを主張します。このテストは、不一致の結合キー、パーティションの回転、またはシフトの回転を失った行などの微妙なバグを捕捉え、またはシフトを繰り返す。

エンドツーエンドパイプラインテスト

エンドツーエンドのテストは、フルライフサイクルをシミュレートします。ソース(例えば、パーケットファイルやカフカトピック)から読み、処理し、ターゲットシンクに書き込みます。これらのテストは外部コンポーネントに依存しているため、専用のテスト環境やコンテナ化されたセットアップ(例えば、Docker Compose with Spark、Mino for object Storage、およびmock Kafka)に最適です。想定されたデータファイルに対する最終出力を検証するか、シンクから読み戻すことで最終テストを検証します。 エンドツーエンドは、毎晩実行されます。

高度なテストの検討

修正を超えて、現代のデータパイプラインは、データ品質、性能SLA、およびレジリエンスを強制しなければなりません。 自動化されたテストは、これらの寸法をカバーすることができます。

データの品質チェックは、Deequでチェック

[Deequ]は、データ品質制約を定義し検証するSparkの上に構築されたライブラリです。 Deequをテストスイートに統合して、完全性(非nullカウント)、ユニークネス(重複するプライマリキーなし)、およびコンプライアンス(例えば、範囲内で落下する値の割合)を確認することができます。制約がテストケースとして各制約を処理します。制約が失敗した場合、対応するテストは失敗します。このアプローチは、最初のレベルのパイプラインが、データが最初に求められているわけではありません。

性能とストレステスト

パイプラインが予想されるデータ量を時間予算内で処理できるかどうか自動パフォーマンステストが測定されます。同じローカルスパークセッションを使用して、テストデータを複数の典型的なバッチサイズにスケールアップします。各ステージの実行期間を記録し、ベースラインとそれを比較します。コード変更が新しいシャッフルまたは非効率的な参加を導入した場合、テストは、再帰を明らかにします。より現実的なパフォーマンスプロファイリングのために、これらのテストを小さなクラスター(例えば、eperald Code を[F] または [Fart] で実行します。[Farget] または [Fart] が、 [Fart] が、 [F] をクラスターに実行します。[F]

CI/CD のテスト

Spark のテストスイートを Jenkins、GitLab CI、または GitHub アクションなどの継続的な統合パイプラインに統合します。パイプラインは以下です。

  • コードをチェックし、テストデータのフィクスチャをロードします。
  • ローカルモードでユニットと統合テストを実行します(高速フィードバック)。
  • パスがすべての場合、オプションで、エンドツーエンドまたはトランスエントクラスターでパフォーマンステストを実行します。
  • テストレポートを発行し、テストが失敗した場合、ビルドに失敗します。

この自動化により、コードがチェックの電池を渡すことなくメインブランチに到達しないことが保証されます。 また、テスト結果の履歴レコードを提供し、特定のコミットへの回帰を追跡するのが容易になります。

メンテナンス可能なテストスイートに最適なプラクティス

  • []Keepテスト独立:[]]])各テストは、独自の入力DataFrameを作成し、共有ミュータブル状態に依存しない。 交差テスト汚染を避けるために、新鮮なスパークセッション(または再利用可能なが、セッションをリセット)を使用してください。
  • [] 代表的なものではなく、小さなデータを使用します。[]] 数ミリ秒で実行するテストは、頻繁に実行を促します。 試験が大きなデータを必要とすると、有意義な結果を生み出すために、一晩実行される遅いCIステージに分けます。
  • [名義のテストは記述的に:[[]のようなテスト名は、読者に、どのような動作が検証されているか、そして期待される結果が何であるかを正確に読者に伝えます。
  • []Refactorテストヘルパー:[ 共通パターンを抽出します(例えば、スパークセッションを作成して、フィクスチャのデータフレームをロード)ユーティリティ機能や特性に。これにより、重複を減らし、パイプラインの変更時にテストスイートを簡単に更新できます。
  • []Version 制御テストデータ:[]] リポジトリにある小さなフィクスチャファイル(例:CSV、パーケット)を格納します。 より大きいデータセットの場合、 ]]DVC などのデータバージョンアップツールを使用して、チェックサム付きのS3バケットに保存します。
  • [:]]] は、パイプラインが、明確なメッセージで例外をスローしたり、適切なときに空のDataFrameを生成したりする、無効な入力を適切に処理していることを確認します。
  • []ドキュメントテストのシナリオ:[]は、各フィクスチャのデータセットの目的と、ビジネスルールがテストされていることを説明するテストディレクトリ内の短いREADMEを維持します。

コンテンツ

Spark ベースのエンジニアリングデータパイプライン用の自動化されたテストフレームワークの構築は、一回限りの努力ではなく、データ信頼性の継続的な投資ではありません。慎重に構築されたテストデータ、よく定義されたアサーション、ローカル実行環境、および CI/CD 統合を組み合わせることで、データエンジニアリングチームはバグを早期にキャッチし、データ品質インシデントを防ぎ、パイプラインの信頼性の変更を出荷することができます。Deequ 制約やパフォーマンスベンチマークなどの高度な技術が、さらに安全網を強化します。その結果は、その決定が最も正確な決定を下すことなく、そのコストを正確に解決するような開発サイクルです。