プロンプトとコンテキスト
アップストリームのジョブは毎日複数のデータアセットを生成します。レポートや品質チェックは、依存関係が更新された後に実行される必要があります。Airflowのアセット認識スケジューリング(asset-aware scheduling)を用いてプロデューサー、コンシューマー、パーティション、リカバリをどのようにモデル化するか、またそれが時間ベースのスケジュールや外部センサーとどのように異なるかを説明してください。
面接官が見ているポイント
- アセットを単なるファイル名やcronラベルではなく、タスクが更新する論理的なデータ依存関係として扱えているか。
- アセットイベント、DAGのtimetable、式の合成、およびイベントトリガーを区別できているか。
- 重複イベント、遅延データ、パーティショングラニュラリティ(粒度)、品質ゲート、およびバックフィルを考慮しているか。
- モニタリング、認可、冪等性、リトライ、および一時停止/再開の動作を説明できるか。
確認すべき明確化のための質問
- アセットはテーブル全体、特定のパーティション、オブジェクトストレージのパスのいずれを表していますか?イベントはパーティションおよびバッチの識別子を保持していますか?
- コンシューマーはすべてのアップストリームアセットを待つ必要がありますか、それともいずれか1つの更新でトリガーできますか?DAG間のアセット式(cross-DAG asset expressions)は必要ですか?
- イベントが到着したものの品質チェックが失敗した場合、コンシューマーをブロックするべきですか、プロデューサーをリトライするべきですか、それとも人間の承認によるリリースが必要ですか?
- 履歴データのバックフィル、イベントの再実行(リプレイ)、または既存のcronスケジュールとの互換性は必要ですか?重複トリガーのコストはどの程度ですか?
30秒での回答例
私はデータプロダクトを更新コントラクトを持つアセットとしてモデル化し、プロデューサーDAGがアトミックな書き込みと品質チェックの成功後にのみアセット更新を発行するようにします。コンシューマーDAGは、より短いcron間隔で準備完了を推測する代わりに、アセットの依存関係や式を使用してトリガー条件を定義します。パーティション、冪等性、重複イベント、遅延更新、バックフィルのルールを定義した上で、イベントのレイテンシ、待機中のDAG、リトライ、鮮度を監視します。更新がAirflowの外部で発生する場合は、イベント駆動トリガーのリカバリ可能性とセキュリティ境界を評価します。
ステップごとの詳細解説
1. アセットコントラクトの確立
Airflow Assetsの公式ドキュメントでは、アセットをDAG間で共有されるデータ依存関係として定義しています。プロデューサーは、データのコミットとコントラクトチェックが成功した後にのみアセットを更新する必要があります。一時ファイルの作成、タスクの開始、またはデータの部分的な書き込みは準備完了を意味しません。コンシューマーが新しいデータと古いデータを区別できるように、アセットURI、オーナー、およびパーティショングラニュラリティを安定した状態に保ちます。
2. トリガーロジックの選択
コンシューマーDAGは1つ以上のアセットに依存できます。公式のAsset-Aware Schedulingドキュメントでは、「すべてのアセットが更新された」または「いずれかのアセットが更新された」などの条件を表現する論理結合がサポートされています。timetableは独立した制約として残すことも可能です。イベント条件と時間条件が一致した場合の動作を定義し、式を行コンテンツのフィルターと誤認しないようにしてください。
3. パーティションと重複の処理
アセットイベントはデータパーティションのウォーターマーク(watermark)を代替するものではありません。追跡可能なバッチまたはパーティションの識別子を保持し、冪等性を確保するためにウォーターマークテーブル、一意キー、またはトランザクション書き込みを使用します。重複イベント、タスクのリトライ、スケジューラのリカバリはいずれも依存関係の再評価を引き起こす可能性があるため、コンシューマーは「イベントが厳密に1回だけ発生する」と仮定せず、安全に再実行可能である必要があります。
4. 外部イベント、品質、およびリカバリ
更新がキューや別のシステムから届く場合は、公式のイベント駆動スケジューリングドキュメントに従ってトリガーを選択し、リカバリ可能なチェック、接続認証情報、およびリソースの解放を検証します。品質チェックが失敗した場合はアセットを未準備状態に保ち、リトライするか明示的なリリースを要求します。バックフィルの場合は、イベントを発行するかどうか、履歴パーティションを分離するかどうか、およびリプレイによってリアルタイムの結果が上書きされないようにするかどうかを決定します。
モデル回答
まず、アセットコントラクト(URI、パーティションキー、プロデューサー、品質ゲート、バッチ識別子)を定義します。プロデューサーDAGは、アトミックな書き込みと品質チェックの後にのみアセットを更新します。コンシューマーDAGは、すべてのアセット更新または任意のアセット更新の条件(必要に応じてtimetableと組み合わせる)にアセット依存関係または式を使用します。イベントにはパーティションとバッチの識別子を持たせ、コンシューマーはウォーターマークと冪等性キーを使用して重複、リトライ、スケジューラのリカバリに対応します。外部更新には、制限された認証情報とリソースライフタイムを持つ制御されたトリガーを使用します。イベントレイテンシ、待機タスク、鮮度、およびリプレイの失敗を監視します。バックフィルとリアルタイムのパスには別々のバッチ境界を使用し、履歴データのリプレイが最新の結果を上書きしないようにします。
よくある間違い
- 更新の所有権や準備完了状態を定義せずに、アセットを単なる任意のパスとして扱う。
- データ準備完了シグナル、パーティションコントラクト、または品質ゲートなしでcron間隔を短縮する。
- アセットイベントが「厳密に1回」配信されると想定し、リトライ、重複、またはスケジューラのリカバリを無視する。
- アセット式を行レベルのコンテンツフィルターとして扱い、実際のパーティションウォーターマークを省略する。
- 外部トリガーがタイムアウトやクリーンアップなしで接続や認証情報を無期限に保持することを許容する。
- バックフィルにリアルタイムDAGを再利用し、履歴のリプレイによって本番の出力が上書きされることを許す。
フォローアップ質問と回答
2つのアップストリームアセットが異なる時間に到着した場合はどうしますか?
コンシューマーがすべてのアセットを待つか、部分的な結果を公開できるかを決定します。すべてのアセットを待つセマンティクスの場合、各パーティションの到着ウォーターマークを追跡し、タイムアウト時にアラートを発報します。部分的な結果の場合はバージョンを公開し、後続のアセットによって更新される可能性があることをコンシューマーに伝えます。
重複イベントをどのようにテストしますか?
テスト環境で同じアセット更新を2回発行し、プロデューサーをリトライさせ、スケジューラを再起動します。行数、バージョン、通知や課金などの外部副作用を含め、コンシューマーの一意キー、ウォーターマーク、トランザクション境界を検証します。
外部センサーを残すのはどのような場合ですか?
依存先システムがAirflowアセットイベントを発行できず、制御されたポーリングインターフェースのみを提供している場合や、移行期間中に互換性を維持する必要がある場合に一時的に残します。移行期限とポーリングコストを記録し、アップストリームシステムを検証可能なアセット更新シグナルへと移行させていきます。