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"
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
)