View a markdown version of this page

Nettoyage de la base de données Aurora PostgreSQL dans un environnement Amazon MWAA - Amazon Managed Workflows for Apache Airflow

Les traductions sont fournies par des outils de traduction automatique. En cas de conflit entre le contenu d'une traduction et celui de la version originale en anglais, la version anglaise prévaudra.

Nettoyage de la base de données Aurora PostgreSQL dans un environnement Amazon MWAA

Amazon Managed Workflows pour Apache Airflow utilise une base de données Aurora PostgreSQL comme base de métadonnées Apache Airflow, où DAG s'exécute et où les instances de tâches sont stockées. L'exemple de code suivant efface régulièrement les entrées de la base de données Aurora PostgreSQL dédiée à votre environnement Amazon MWAA.

Important

Apache Airflow v3 restreint l'accès direct à la base de données de métadonnées à partir du code de tâche. Les travailleurs ne se connectent plus à la base de données de métadonnées et le DAG ou le code de tâche ne peuvent pas importer ou utiliser directement les sessions ou les modèles de base de données Apache Airflow. Ce changement améliore la sécurité et l'évolutivité. Cependant, l'approche DAG-based de nettoyage de base de données qui fonctionne dans Apache Airflow v2 ne fonctionne pas dans les environnements Apache Airflow v3.

Utilisez plutôt la commande airflow db clean CLI via le point de terminaison de la CLI Amazon MWAA pour effectuer le nettoyage de la base de données de métadonnées.

Note

Au fil du temps, la base de données de métadonnées accumule d'anciens enregistrements, des données XCom et des données de tâches périmées. Cette croissance utilise les connexions aux bases de données, ralentit votre environnement et retarde les tâches. Effectuez un nettoyage régulier des métadonnées pour éviter ces problèmes.

Version

Les exemples de code de cette page sont spécifiques à Apache Airflow v2 et v3 pris en charge sur Amazon MWAA. Reportez-vous aux versions d'Apache Airflow prises en charge.

Conditions préalables

Pour utiliser l'exemple de code de cette page, vous aurez besoin des éléments suivants :

Dépendances

Pour utiliser cet exemple de code avec Apache Airflow v2, aucune dépendance supplémentaire n'est requise. Utilisez https://github.com/aws/amazon-mwaa-docker-images aws-mwaa-docker-images pour installer Apache Airflow.

Exemple de code

Les exemples suivants montrent comment nettoyer la base de données de métadonnées de votre environnement Amazon MWAA.

Apache Airflow v3.0.6 to 3.2.1
Considérations importantes
  • Vous devez spécifier --clean-before-timestamp pour contrôler la distance parcourue par le nettoyage. Utilisez un horodatage au format ISO 8601 (par exemple,). 2025-01-01T00:00:00+00:00

  • Nous vous recommandons de spécifier --tables de limiter le nettoyage à des tables spécifiques. En cas d'omission, la commande nettoie toutes les tables prises en charge.

  • Commencez par une petite portée : utilisez d'abord une --clean-before-timestamp valeur plus ancienne (plus proche de la date de création de votre environnement) et un seul tableau. Cela limite le nettoyage aux seuls enregistrements les plus anciens. Étant donné que la commande supprime tout ce qui se trouve avant l'horodatage spécifié, l'utilisation d'un horodatage plus récent entraîne une plus grande étendue de suppression. Avancez progressivement l'horodatage au fur et à mesure que vous gagnez en confiance dans le processus.

  • Large-scale le nettoyage peut avoir un impact sur les performances de la base de données. La suppression d'un volume élevé d'enregistrements exerce une pression sur la base de données Aurora PostgreSQL et peut affecter la réactivité de votre environnement. Utilisez le --batch-size paramètre pour contrôler la taille des transactions et envisagez d'exécuter le nettoyage pendant les périodes de faible trafic. Faites preuve de prudence lors de l'exécution dans des environnements de production.

Comportement des tables dépendantes

Lorsque vous spécifiez une table avec--tables, la commande inclut automatiquement toutes les tables dépendantes (enfants) qui ont des relations de clé étrangère avec la table spécifiée. Les enregistrements de la table enfant sont d'abord supprimés, puis les enregistrements de la table parent, pour satisfaire aux contraintes de clé étrangère. Par exemple, la spécification --tables dag_run permet également de nettoyertask_instance, task_instance_history xcomtask_state_store, et deadline parce que ces tables font référence dag_run à des clés étrangères.

Le tableau suivant récapitule les chaînes de dépendance.

Tableau spécifié Tables supplémentaires nettoyées (personnes à charge)
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

Les tables sans dépendances (telles quelog, jobimport_error,sla_miss) sont nettoyées de manière isolée lorsque cela est spécifié.

--dry-runÀ utiliser pour voir exactement quelles tables et combien de lignes seraient affectées avant de procéder à un nettoyage.

Paramètres disponibles pour airflow db clean

Le tableau suivant décrit les paramètres disponibles pourairflow db clean.

Paramètre Description Par défaut
--clean-before-timestamp (Obligatoire) Date ou horodatage avant laquelle les données sont purgées. Si aucun fuseau horaire n'est fourni, le fuseau horaire par défaut d'Apache Airflow est utilisé. Exemple : 2025-01-01T00:00:00+00:00 Aucune
--tables ou -t Noms des tables sur lesquelles effectuer la maintenance (séparés par des virgules). Les options sont les suivantes : dag_run task_instance task_instance_history log job xcomimport_error,task_reschedule,trigger, dagdag_version,sla_miss,callback_request, celery_taskmetacelery_tasksetmeta,asset_event,deadline, revoked_tokentask_state_store,connection_test_request, _xcom_archive Aucune
--batch-size Nombre maximum de lignes à supprimer ou à archiver en une seule transaction. Des valeurs plus faibles réduisent les blocages de longue durée mais augmentent le nombre de lots. Aucune
--dry-run Effectuez un essai à sec sans supprimer réellement les données. Recommandé pour les tests initiaux. False
--skip-archive Ne conservez pas les enregistrements purgés dans une table d'archives. Par défaut, db clean déplace les enregistrements purgés dans des tables d'archives (nommées selon une _<table>_archive convention, par exemple_dag_run_archive) au lieu de les supprimer définitivement. Cela constitue un filet de sécurité : vous pouvez inspecter les données archivéesairflow db export-archived, les exporter avec ou les supprimer ultérieurementairflow db drop-archived. Lorsque cette --skip-archive option est définie, les enregistrements sont définitivement supprimés sans cette étape intermédiaire. False
--dag-ids Seules les données de nettoyage liées aux identifiants DAG donnés. Aucune
--exclude-dag-ids Évitez de nettoyer les données relatives aux identifiants DAG donnés. Aucune
-y, --yes Ignorez l'invite de confirmation. Requis pour l'exécution non interactive de la CLI via Amazon MWAA. False
-v, --verbose Rendez la sortie de journalisation plus détaillée. False

Pour plus d'informations sur les paramètres disponibles, consultez la référence des variables CLI et env sur le site Web d'Apache Airflow.

Exemples de code

Les exemples suivants montrent comment appeler airflow db clean via le point de terminaison de la CLI Amazon MWAA. Pour plus d'informations sur la création de jetons CLI, consultezCréation d'un jeton CLI Apache Airflow.

À l'aide d'un script 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}")
Exemple de fonctionnement à sec (première étape recommandée)

Avant d'effectuer un véritable nettoyage, exécutez with --dry-run pour voir ce qui serait supprimé :

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 )