翻訳は機械翻訳により提供されています。提供された翻訳内容と英語版の間で齟齬、不一致または矛盾がある場合、英語版が優先します。
Amazon MWAA 環境での Aurora PostgreSQL データベースのクリーンアップ
Amazon Managed Workflows for Apache Airflow は、Apache Airflow メタデータデータベースとして Aurora PostgreSQL データベースを使用します。このメタデータデータベースには DAG が実行され、タスクインスタンスが保存されます。次のサンプルコードは、Amazon MWAA 環境の専用の Aurora PostgreSQL データベースから定期的にエントリを消去します。
Apache Airflow v3 は、タスクコードからのメタデータデータベースへの直接アクセスを制限します。ワーカーはメタデータデータベースに接続しなくなり、DAG またはタスクコードは Apache Airflow データベースセッションまたはモデルを直接インポートまたは使用することはできません。この変更により、セキュリティとスケーラビリティが向上します。ただし、Apache Airflow v2 で動作する DAG ベースのデータベースクリーンアップアプローチは、Apache Airflow v3 環境では機能しません。
代わりに、Amazon airflow db clean MWAA CLI エンドポイントを介して CLI コマンドを使用して、メタデータデータベースのクリーンアップを実行します。
時間の経過とともに、メタデータデータベースは古いレコード、XCom データ、古いタスクデータを蓄積します。この増加により、データベース接続が使用され、環境が遅くなり、タスクが遅延します。これらの問題を防ぐために、定期的なメタデータクリーンアップを実行します。
バージョン
このページのコードサンプルは、Amazon MWAA でサポートされている Apache Airflow v2 および v3 に固有のものです。サポートされている Apache Airflow バージョン を参照してください。
前提条件
このページのサンプルコードを使用するには、以下が必要です。
依存関係
このコード例を Apache Airflow v2 で使用する場合、追加の依存関係は必要ありません。aws-mwaa-docker-images を使用して、Apache Airflow をインストールします。
コードサンプル
次の例は、Amazon MWAA 環境でメタデータデータベースをクリーンアップする方法を示しています。
- Apache Airflow v3.0.6 to 3.2.1
-
重要な考慮事項
-
クリーンアップの到達距離を制御する--clean-before-timestampには、 を指定する必要があります。ISO 8601 形式のタイムスタンプ ( など2025-01-01T00:00:00+00:00) を使用します。
-
クリーンアップを特定のテーブルに制限--tablesするには、 を指定することをお勧めします。省略すると、コマンドはサポートされているすべてのテーブルをクリーンアップします。
-
小さなスコープから始める — まず古い--clean-before-timestamp値 (環境作成日に近い) と 1 つのテーブルを使用します。これにより、クリーンアップは最も古いレコードのみに制限されます。コマンドは指定されたタイムスタンプより前にすべてを削除するため、最新のタイムスタンプを使用すると、削除範囲が大きくなります。プロセスに自信が持てたら、タイムスタンプを徐々に前方に移動します。
-
大規模なクリーンアップは、データベースのパフォーマンスに影響を与える可能性があります。大量のレコードを削除すると、Aurora PostgreSQL データベースに負荷がかかり、環境の応答性に影響する可能性があります。--batch-size パラメータを使用してトランザクションサイズを制御し、トラフィックが少ない時間帯にクリーンアップを実行することを検討してください。本番環境で を実行するときは注意が必要です。
依存テーブルの動作
でテーブルを指定すると--tables、コマンドには、指定されたテーブルと外部キー関係を持つ依存 (子) テーブルが自動的に含まれます。外部キーの制約を満たすために、子テーブルレコードが最初に削除され、次に親テーブルレコードが削除されます。例えば、 を指定するとtask_instance、これらのテーブルは外部キーdag_runを介して参照deadlineされるため、task_state_store、、task_instance_historyxcom、、および --tables dag_runもクリーンアップされます。
次の表は、依存関係チェーンをまとめたものです。
| 指定されたテーブル |
追加テーブルのクリーンアップ (依存関係) |
dag_run |
task_instance, task_instance_history, xcom, task_state_store, deadline |
dag |
dag_version, deadline |
task_instance |
task_instance_history, xcom |
trigger |
task_instance, task_instance_history, xcom |
dag_version |
task_instance, task_instance_history, xcom, dag_run |
依存関係のないテーブル (log、job、import_error、 などsla_miss) は、指定されたときに個別にクリーンアップされます。
--dry-run を使用して、クリーンアップにコミットする前に影響を受けるテーブルと行数を正確に確認します。
で使用できるパラメータ airflow db clean
次の表に、 で使用できるパラメータを示しますairflow db clean。
| パラメータ |
説明 |
デフォルト |
--clean-before-timestamp |
(必須) データが消去される前の日付またはタイムスタンプ。タイムゾーンが指定されていない場合、Apache Airflow のデフォルトのタイムゾーンが想定されます。例: 2025-01-01T00:00:00+00:00 |
なし |
--tables、または -t |
メンテナンスを実行するテーブル名 (カンマ区切り)。オプションにはdag_run、、task_instance、task_instance_history、、log、jobxcom、import_error、、task_rescheduletrigger、、dag、dag_version、、sla_miss、callback_request、celery_taskmeta、、celery_tasksetmeta、asset_event、connection_test_request、、、、、 deadline revoked_token task_state_storeなどがあります。 _xcom_archive |
なし |
--batch-size |
1 つのトランザクションで削除またはアーカイブする行の最大数。値を小さくすると、長時間実行されるロックは減少しますが、バッチ数は増加します。 |
なし |
--dry-run |
実際にデータを削除せずにドライランを実行します。初期テストに推奨されます。 |
誤 |
--skip-archive |
消去されたレコードをアーカイブテーブルに保持しないでください。デフォルトでは、 は消去されたレコードを完全に削除する_dag_run_archiveのではなく、アーカイブテーブル ( などの_<table>_archive規則で命名) db cleanに移動します。これにより、アーカイブされたデータを検査したり、 でエクスポートしたりairflow db export-archived、後で で削除したりできますairflow db drop-archived。を設定する--skip-archiveと、レコードはこの中間ステップなしで完全に削除されます。 |
誤 |
--dag-ids |
指定された DAG IDs。 |
なし |
--exclude-dag-ids |
指定された DAG IDs。 |
なし |
-y, --yes |
確認プロンプトをスキップします。Amazon MWAA を介した非インタラクティブ CLI の実行に必要です。 |
誤 |
-v, --verbose |
ログ記録出力をより詳細にします。 |
誤 |
使用可能なパラメータの詳細については、Apache Airflow ウェブサイトの「CLI および env 変数リファレンス」を参照してください。
コードサンプル
次の例は、Amazon MWAA CLI エンドポイントairflow db cleanを介して を呼び出す方法を示しています。CLI トークンの作成の詳細については、「」を参照してくださいApache Airflow CLI トークンの作成。
Python スクリプトの使用:
import boto3
import base64
import requests
# Replace with your environment name and AWS Region
mwaa_env_name = "YOUR_ENVIRONMENT_NAME"
region = "YOUR_REGION"
# Configure cleanup scope
clean_before_timestamp = "2025-06-01T00:00:00+00:00"
tables = "dag_run,task_instance,log,job,xcom"
# Build the Airflow CLI command
# -y flag is required to skip interactive confirmation prompt
airflow_cmd = f"db clean --clean-before-timestamp {clean_before_timestamp} --tables {tables} -y"
# Create a CLI token
client = boto3.client("mwaa", region_name=region)
cli_token_response = client.create_cli_token(Name=mwaa_env_name)
cli_token = cli_token_response["CliToken"]
web_server_hostname = cli_token_response["WebServerHostname"]
# Invoke the Airflow CLI through the MWAA endpoint
url = f"https://{web_server_hostname}/aws_mwaa/cli"
response = requests.post(
url,
headers={
"Authorization": f"Bearer {cli_token}",
"Content-Type": "text/plain",
},
data=airflow_cmd,
)
# Parse and display the results
stdout_message = base64.b64decode(response.json()["stdout"]).decode("utf-8")
stderr_message = base64.b64decode(response.json()["stderr"]).decode("utf-8")
print(f"Status code: {response.status_code}")
print(f"stdout:\n{stdout_message}")
print(f"stderr:\n{stderr_message}")
ドライランの例 (推奨される最初のステップ)
実際のクリーンアップを実行する前に、 --dry-run で を実行して、削除される内容を確認します。
AIRFLOW_CMD="db clean --clean-before-timestamp 2025-06-01T00:00:00+00:00 --tables dag_run,task_instance --dry-run -y"
- Apache Airflow v2.7.2 to 2.11.2
-
from airflow import DAG
from airflow.models.param import Param
from airflow.operators.bash_operator import BashOperator
from airflow.utils.dates import days_ago
from datetime import datetime, timedelta
# Note: Database commands might time out if running longer than 5 minutes. If this occurs, please increase the MAX_AGE_IN_DAYS (or change
# timestamp parameter to an earlier date) for initial runs, then reduce on subsequent runs until the desired retention is met.
MAX_AGE_IN_DAYS = 30
# To clean specific tables, please provide a comma-separated list per
# https://airflow.apache.org/docs/apache-airflow/stable/cli-and-env-variables-ref.html#clean
# A value of None will clean all tables
TABLES_TO_CLEAN = None
with DAG(
dag_id="clean_db_dag",
schedule_interval=None,
catchup=False,
start_date=days_ago(1),
params={
"timestamp": Param(
default=(datetime.now()-timedelta(days=MAX_AGE_IN_DAYS)).strftime("%Y-%m-%d %H:%M:%S"),
type="string",
minLength=1,
maxLength=255,
),
}
) as dag:
if TABLES_TO_CLEAN:
bash_command="airflow db clean --clean-before-timestamp '{{ params.timestamp }}' --tables '"+TABLES_TO_CLEAN+"' --skip-archive --yes"
else:
bash_command="airflow db clean --clean-before-timestamp '{{ params.timestamp }}' --skip-archive --yes"
cli_command = BashOperator(
task_id="bash_command",
bash_command=bash_command
)