View a markdown version of this page

Amazon MWAA での Apache Airflow のパフォーマンス調整 - Amazon Managed Workflows for Apache Airflow

翻訳は機械翻訳により提供されています。提供された翻訳内容と英語版の間で齟齬、不一致または矛盾がある場合、英語版が優先します。

Amazon MWAA での Apache Airflow のパフォーマンス調整

このトピックでは、Amazon MWAA での Apache Airflow 設定オプションの使用 を使用して Amazon Managed Workflows for Apache Airflow 環境のパフォーマンスを調整する方法について説明します。

Apache Airflow 設定オプションの追加

環境に Airflow 設定オプションを追加するには、次の手順に従います。

  1. Amazon MWAA コンソールで、環境ページ を開きます。

  2. 環境を選択します。

  3. 編集 を選択します。

  4. 次へ を選択します。

  5. Airflow 設定オプション ペインで カスタム設定を追加 を選択します。

  6. ドロップダウンリストから設定を選択して値を入力するか、カスタム設定を入力して値を入力します。

  7. 追加する設定ごとに カスタム設定を追加 を選択します。

  8. 保存 を選択します。

詳細については、Amazon MWAA での Apache Airflow 設定オプションの使用 を参照してください。

Apache Airflow スケジューラー

Apache Airflow スケジューラーは Apache Airflow のコアコンポーネントです。スケジューラーに問題があると、DAG の解析やタスクのスケジュール設定ができなくなる可能性があります。Apache Airflow スケジューラーの調整の詳細については、Apache Airflow ドキュメンテーションウェブサイトの スケジューラーのパフォーマンスのファインチューニング を参照してください。

パラメータ

このセクションでは、Apache Airflow スケジューラー (Apache Airflow v2 以降) で使用できる設定オプションとそのユースケースについて説明します。

Apache Airflow v3
設定 ユースケース

celery.sync_parallelism

Celery エグゼキュターがタスクの状態を同期するために使用するプロセスの数です。

デフォルト: 1

このオプションを使用して、Celery Executor が使用するプロセスを制限することで、キューの競合を防止します。デフォルトでは、CloudWatch Logs にタスクログを配信する際にエラーが発生しないように値が 1 に設定されています。この値を 0 に設定すると最大数のプロセスを使用することになりますが、タスクログの配信時にエラーが発生する可能性があります。

scheduler.scheduler_idle_sleep_time

ループに何もしなければ、スケジューラの「ループ」が連続して待機する秒数。

デフォルト: 1

このオプションを使用すると、「ループ」の完了後にスケジューラがスリープする時間を増やすことで、スケジューラの CPU 使用率を解放できます。この値を大きくすると、Apache Airflow v2 および Apache Airflow v3 dag_processor.parsing_processesの で使用可能なスケジューラスレッドが減少します。これにより、スケジューラーが DAG を解析する容量が減り、DAG がウェブサーバーに取り込まれるのにかかる時間が長くなる可能性があります。

scheduler.max_dagruns_to_create_per_loop

スケジューラー「ループ」ごとに作成する DagRuns の DAG の最大数。

デフォルト: 10

このオプションを使用すると、スケジューラの「ループ」の DagRuns の最大数を減らすことで、タスクをスケジュールするためのリソースを解放できます。

dag_processor.parsing_processes

スケジューラーが DAG をスケジュールするために並列して実行できるスレッドの数。

デフォルト: (2 * number of vCPUs) - 1 を使用する

このオプションを使用すると、スケジューラが DAG を解析するために並行して実行するプロセスの数を減らすことで、リソースを解放できます。 DAGs DAG の解析がタスクスケジューリングに影響している場合は、この数を低く抑えることをお勧めします。環境の vCPU 数よりも小さい値を指定する 必要があります。詳細については、制限 を参照してください。

Apache Airflow v2
設定 ユースケース

celery.sync_parallelism

Celery エグゼキュターがタスクの状態を同期するために使用するプロセスの数です。

デフォルト: 1

このオプションを使用して、Celery Executor が使用するプロセスを制限することで、キューの競合を防止します。デフォルトでは、CloudWatch Logs にタスクログを配信する際にエラーが発生しないように値が 1 に設定されています。この値を 0 に設定すると最大数のプロセスを使用することになりますが、タスクログの配信時にエラーが発生する可能性があります。

scheduler.scheduler_idle_sleep_time

ループに何もしなければ、スケジューラの「ループ」が連続して待機する秒数。

デフォルト: 1

このオプションを使用すると、「ループ」の完了後にスケジューラがスリープする時間を増やすことで、スケジューラの CPU 使用率を解放できます。この値を大きくすると、Apache Airflow v2 および Apache Airflow v3 dag_processor.parsing_processesの で使用可能なスケジューラスレッドが減少します。これにより、スケジューラーが DAG を解析する容量が減り、DAG がウェブサーバーに取り込まれるのにかかる時間が長くなる可能性があります。

scheduler.max_dagruns_to_create_per_loop

スケジューラー「ループ」ごとに作成する DagRuns の DAG の最大数。

デフォルト: 10

このオプションを使用すると、スケジューラの「ループ」の DagRuns の最大数を減らすことで、タスクをスケジュールするためのリソースを解放できます。

scheduler.parsing_processes

スケジューラーが DAG をスケジュールするために並列して実行できるスレッドの数。

デフォルト: (2 * number of vCPUs) - 1 を使用する

このオプションを使用すると、スケジューラが DAG を解析するために並行して実行するプロセスの数を減らすことで、リソースを解放できます。 DAGs DAG の解析がタスクスケジューリングに影響している場合は、この数を低く抑えることをお勧めします。環境の vCPU 数よりも小さい値を指定する 必要があります。詳細については、制限 を参照してください。

制限

このセクションでは、スケジューラーのデフォルトパラメータを調整する際に考慮する制限について説明します。

scheduler.parsing_processes、scheduler.max_threads (v2のみ)

1 つの環境クラスの vCPU ごとに 2 つのスレッドを使用できます。環境クラスのスケジューラー用に少なくとも 1 つのスレッドを予約する必要があります。タスクのスケジュールが遅れていることに気付いた場合は、環境クラス を増やす必要があるかもしれません。例えば、大規模な環境では、スケジューラ用に 4 vCPU の Fargate コンテナインスタンスがあります。つまり、他のプロセスに使用できるスレッドの 7 合計は最大数になります。つまり、2 つのスレッドに 4 つの vCPUs を乗算した値から、スケジューラー自体の 1 を引いた値です。scheduler.max_threads (v2のみ) および scheduler.parsing_processes で指定する値は、環境クラスで使用できるスレッド数 (以下を参照) を超えてはなりません。

  • mw1.small — 他のプロセスの 1 スレッド数を超えてはいけません。残りのスレッドはスケジューラー用に予約されています。

  • mw1.medium — 他のプロセスの 3 スレッド数を超えてはいけません。残りのスレッドはスケジューラー用に予約されています。

  • mw1.large — 他のプロセスの 7 スレッド数を超えてはいけません。残りのスレッドはスケジューラー用に予約されています。

DAG フォルダー

Apache Airflow スケジューラーは、環境内の DAG フォルダーを継続的にスキャンします。含まれている plugins.zip ファイル、または「airflow」インポートステートメントを含む Python (.py) ファイル。生成されたすべての Python DAG オブジェクトは DagBag に格納され、そのファイルがスケジューラーによって処理され、スケジュールする必要があるタスク (ある場合) が決定されます。DAG ファイルの解析は、そのファイルに実行可能な DAG オブジェクトが含まれているかどうかに関係なく行われます。

パラメータ

このセクションでは、DAG フォルダー (Apache Airflow v2 以降) で使用できる設定オプションとそのユースケースについて説明します。

Apache Airflow v3
設定 ユースケース

dag_processor.refresh_interval

DAG フォルダーをスキャンして新しいファイルを探す必要のある秒数。

デフォルト: 300 秒

このオプションを使用して、DAGsフォルダを解析する秒数を増やしてリソースを解放します。DAG フォルダーに大量のファイルがあることが原因で、total_parse_time metrics で解析時間が長くなる場合は、この値を増やすことをお勧めします。

dag_processor.min_file_process_interval

スケジューラーが DAG を解析して DAG への更新が反映されるまでの秒数。

デフォルト: 30 秒

このオプションを使用して、DAG を解析する前にスケジューラが待機する秒数を増やしてリソースを解放します。例えば、30 の値を指定すると、DAG ファイルは 30 秒ごとに解析されます。お使いの環境の CPU 使用率を下げるには、この数値を高くしておくことをお勧めします。

Apache Airflow v2
設定 ユースケース

scheduler.dag_dir_list_interval

DAG フォルダーをスキャンして新しいファイルを探す必要のある秒数。

デフォルト: 300 秒

このオプションを使用して、DAGsフォルダを解析する秒数を増やしてリソースを解放します。DAG フォルダーに大量のファイルがあることが原因で、total_parse_time metrics で解析時間が長くなる場合は、この値を増やすことをお勧めします。

scheduler.min_file_process_interval

スケジューラーが DAG を解析して DAG への更新が反映されるまでの秒数。

デフォルト: 30 秒

このオプションを使用して、DAG を解析する前にスケジューラが待機する秒数を増やしてリソースを解放します。例えば、30 の値を指定すると、DAG ファイルは 30 秒ごとに解析されます。お使いの環境の CPU 使用率を下げるには、この数値を高くしておくことをお勧めします。

DAG ファイル

Apache Airflow スケジューラーループの一部として、個々の DAG ファイルが解析されて DAG Python オブジェクトが抽出されます。Apache Airflow v2 以降では、スケジューラーは同時に最大数の 解析プロセス を解析します。同じファイルが再び解析される前に、scheduler.min_file_process_interval (v2) または dag_processor.min_file_process_interval (v3) で指定された秒数が経過する必要があります。

パラメータ

このセクションでは、Apache Airflow DAG ファイル (Apache Airflow v2 以降) で使用できる設定オプションとそのユースケースについて説明します。

Apache Airflow v3
設定 ユースケース

dag_processor.dag_file_processor_timeout

DagFileProcessor が DAG ファイルの処理をタイムアウトするまでの秒数。

デフォルト: 50 秒

このオプションを使用すると、DagFileProcessor がタイムアウトするまでにかかる時間が長くなります。DAG 処理ログにタイムアウトが発生し、実行可能な DAG が読み込まれなくなる場合は、この値を増やすことをお勧めします。

core.dagbag_import_timeout

Python ファイルをインポートするまでの秒数がタイムアウトします。

デフォルト: 30 秒

このオプションを使用すると、Python ファイルのインポート中にスケジューラがタイムアウトして DAG オブジェクトを抽出するまでにかかる時間が長くなります。このオプションはスケジューラーの「ループ」の一部として処理され、dag_processor.dag_file_processor_timeout で指定されている値よりも小さい値が含まれている必要があります。

core.min_serialized_dag_update_interval

データベース内のシリアル化された DAG が更新されるまでの最小秒数。

デフォルト: 30

このオプションを使用して、データベース内のシリアル化された DAGs が更新されるまでの秒数を増やしてリソースを解放します。DAG の数が多い場合、または複雑な DAG がある場合は、この値を増やすことをおすすめします。この値を増やすと、DAG がシリアル化されるときのスケジューラーとデータベースへの負荷が軽減されます。

core.min_serialized_dag_fetch_interval

シリアル化された DAG が DAGBag に既に読み込まれているときに、データベースから再フェッチされる秒数。

デフォルト: 10

このオプションを使用して、シリアル化された DAG が再取得される秒数を増やしてリソースを解放します。データベースの「書き込み」速度を下げるには、この値を core.min_serialized_dag_update_interval で指定した値より大きくする必要があります。この値を増やすと、DAG がシリアル化される際のウェブサーバーとデータベースの負荷が軽減されます。

Apache Airflow v2
設定 ユースケース

core.dag_file_processor_timeout

DagFileProcessor が DAG ファイルの処理をタイムアウトするまでの秒数。

デフォルト: 50 秒

このオプションを使用すると、DagFileProcessor がタイムアウトするまでにかかる時間が長くなります。DAG 処理ログにタイムアウトが発生し、実行可能な DAG が読み込まれなくなる場合は、この値を増やすことをお勧めします。

core.dagbag_import_timeout

Python ファイルをインポートするまでの秒数がタイムアウトします。

デフォルト: 30 秒

このオプションを使用すると、Python ファイルのインポート中にスケジューラがタイムアウトして DAG オブジェクトを抽出するまでにかかる時間が長くなります。このオプションはスケジューラーの「ループ」の一部として処理され、core.dag_file_processor_timeout で指定されている値よりも小さい値が含まれている必要があります。

core.min_serialized_dag_update_interval

データベース内のシリアル化された DAG が更新されるまでの最小秒数。

デフォルト: 30

このオプションを使用して、データベース内のシリアル化された DAGs が更新されるまでの秒数を増やしてリソースを解放します。DAG の数が多い場合、または複雑な DAG がある場合は、この値を増やすことをおすすめします。この値を増やすと、DAG がシリアル化されるときのスケジューラーとデータベースへの負荷が軽減されます。

core.min_serialized_dag_fetch_interval

シリアル化された DAG が DAGBag に既に読み込まれているときに、データベースから再フェッチされる秒数。

デフォルト: 10

このオプションを使用して、シリアル化された DAG が再取得される秒数を増やしてリソースを解放します。データベースの「書き込み」速度を下げるには、この値を core.min_serialized_dag_update_interval で指定した値より大きくする必要があります。この値を増やすと、DAG がシリアル化される際のウェブサーバーとデータベースの負荷が軽減されます。

タスク

Apache Airflow スケジューラーとワーカーはどちらもタスクのキューイングとデキューに関与します。スケジューラーは、解析済みのスケジュール設定が完了したタスクを なし ステータスから スケジュール済み ステータスに移行します。同じく Fargate のスケジューラーコンテナーで実行されているエグゼキューターは、それらのタスクをキューに入れ、そのステータスを キューで待機中 に設定します。ワーカーにキャパシティがあると、キューからタスクを取り出してステータスを 実行中 に設定し、その後、タスクが成功したか失敗したかに基づいてステータスを 成功 または 失敗 に変更します。

パラメータ

このセクションでは、Apache Airflow タスクで使用できる設定オプションとそのユースケースについて説明します。

Amazon MWAA がオーバーライドするデフォルトの設定オプションは でマークされています。

Apache Airflow v3
設定 ユースケース

core.parallelism

各スケジューラが同時にモニタリングおよび実行できるタスクインスタンスの最大数。

デフォルト: (maxWorkers * maxCeleryWorkers) / schedulers * 1.5 に基づいて動的に設定されます。

このオプションを使用して、実行中のタスクの合計のハードグローバル上限を制御します。例えば、この値を大きくすると、メタデータベースを同時接続が多すぎたり、コストを抑えたりすることを防ぐことができます。デフォルト値は高いため、他のコントロール (worker_autoscale、プールスロット、max_active_tasks_per_dag) が有効な同時実行制限として機能します。

core.max_active_tasks_per_dag

各 DAG 実行で同時に実行できるタスクインスタンスの最大数。

デフォルト: 16

このオプションを使用すると、同時に実行できるタスクインスタンスの数を増やしてリソースを解放できます。たとえば、10 個の並列タスクを持つ DAGs があり、すべての DAGs を同時に実行する場合は、最大並列度を計算します。使用可能なワーカーの数に のタスク密度を乗算しcelery.worker_concurrency、DAGs。

core.execute_tasks_new_python_interpreter

Apache Airflow が親プロセスをフォークするか、新しい Python プロセスを作成してタスクを実行するかを決定します。

デフォルト: True

True に設定すると、Apache Airflow はプラグインに加えた変更を、タスクを実行するために作成された新しい Python プロセスとして認識します。

celery.worker_concurrency

Amazon MWAA は、このオプションの Airflow ベースインストールをオーバーライドして、自動スケーリングコンポーネントの一部としてワーカーをスケーリングします。

デフォルト: 該当なし

このオプションに指定された値はすべて無視されます。

celery.worker_autoscale

ワーカー向けのタスク同時実行性。

デフォルト:

  • mw1.micro - 3,0

  • mw1.small - 5,0

  • mw1.medium - 10,0

  • mw1.large - 20,0

  • mw1.xlarge - 40,0

  • mw1.2xlarge - 80,0

このオプションを使用すると、ワーカーの maximumminimumタスク同時実行を減らすことでリソースを解放できます。ワーカーは、そのための十分なリソースがあるかどうかに関係なく、構成された maximum までの同時タスク受け入れます。十分なリソースがない状態でタスクをスケジュールすると、タスクはすぐに失敗します。リソースを大量に消費するタスクの場合は、タスクあたりの容量を増やすために値をデフォルトより小さくしてこの値を変更することをお勧めします。

Apache Airflow v2
設定 ユースケース

core.parallelism

各スケジューラが同時にモニタリングおよび実行できるタスクインスタンスの最大数。

デフォルト: (maxWorkers * maxCeleryWorkers) / schedulers * 1.5 に基づいて動的に設定されます。

このオプションを使用して、実行中のタスクの合計のハードグローバル上限を制御します。例えば、この値を大きくすると、メタデータベースを同時接続が多すぎたり、コストを抑えたりすることを防ぐことができます。デフォルト値は高いため、他のコントロール (worker_autoscale、プールスロット、max_active_tasks_per_dag) が有効な同時実行制限として機能します。

core.dag_concurrency

各 DAG で同時に実行できるタスクインスタンスの数。

デフォルト: 16

この設定は Airflow 2.2.0 以降廃止され、 に置き換えられましたcore.max_active_tasks_per_dag

このオプションを使用すると、同時に実行できるタスクインスタンスの数を増やしてリソースを解放できます。たとえば、10 個の並列タスクを持つ DAGs があり、すべての DAGs を同時に実行する場合は、最大並列度を計算します。使用可能なワーカーの数に のタスク密度を乗算しcelery.worker_concurrency、DAGs。

core.execute_tasks_new_python_interpreter

Apache Airflow が親プロセスをフォークするか、新しい Python プロセスを作成してタスクを実行するかを決定します。

デフォルト: True

True に設定すると、Apache Airflow はプラグインに加えた変更を、タスクを実行するために作成された新しい Python プロセスとして認識します。

celery.worker_concurrency

Amazon MWAA は、このオプションの Airflow ベースインストールをオーバーライドして、自動スケーリングコンポーネントの一部としてワーカーをスケーリングします。

デフォルト: 該当なし

このオプションに指定された値はすべて無視されます。

celery.worker_autoscale

ワーカー向けのタスク同時実行性。

デフォルト:

  • mw1.micro - 3,0

  • mw1.small - 5,0

  • mw1.medium - 10,0

  • mw1.large - 20,0

  • mw1.xlarge - 40,0

  • mw1.2xlarge - 80,0

このオプションを使用すると、ワーカーの maximumminimumタスク同時実行を減らすことでリソースを解放できます。ワーカーは、そのための十分なリソースがあるかどうかに関係なく、構成された maximum までの同時タスク受け入れます。十分なリソースがない状態でタスクをスケジュールすると、タスクはすぐに失敗します。リソースを大量に消費するタスクの場合は、タスクあたりの容量を増やすために値をデフォルトより小さくしてこの値を変更することをお勧めします。