Considerazioni importanti
-
È necessario specificare --clean-before-timestamp per controllare la distanza percorsa dalla pulizia. Utilizza un timestamp in formato ISO 8601 (ad esempio,). 2025-01-01T00:00:00+00:00
-
Ti consigliamo di specificare di limitare --tables la pulizia a tabelle specifiche. Se omesso, il comando pulisce tutte le tabelle supportate.
-
Inizia con un ambito piccolo: utilizza prima un --clean-before-timestamp valore precedente (più vicino alla data di creazione dell'ambiente) e una singola tabella. Ciò limita la pulizia solo ai record più vecchi. Poiché il comando elimina tutto ciò che precede il timestamp specificato, l'utilizzo di un timestamp più recente comporta un ambito di eliminazione più ampio. Sposta gradualmente il timestamp in avanti man mano che acquisisci sicurezza nel processo.
-
Large-scale la pulizia può influire sulle prestazioni del database. L'eliminazione di un volume elevato di record esercita una pressione sul database Aurora PostgreSQL e potrebbe influire sulla reattività dell'ambiente. Utilizzate il --batch-size parametro per controllare le dimensioni delle transazioni e valutate la possibilità di eseguire Cleanup nei periodi di traffico ridotto. Fai attenzione quando esegui l'esecuzione in ambienti di produzione.
Comportamento dipendente della tabella
Quando si specifica una tabella con--tables, il comando include automaticamente tutte le tabelle dipendenti (secondarie) che hanno relazioni di chiave esterna con la tabella specificata. I record della tabella secondaria vengono eliminati prima, quindi i record della tabella principale, per soddisfare i vincoli di chiave esterna. Ad esempio, specificando si eliminano --tables dag_run anche i task_instance valori, task_instance_history xcomtask_state_store, e deadline perché queste tabelle fanno riferimento dag_run tramite chiavi esterne.
La tabella seguente riassume le catene di dipendenze.
| Tabella specificata |
Tabelle aggiuntive pulite (dipendenti) |
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 |
Le tabelle senza elementi dipendenti (ad esempiolog,, jobimport_error,sla_miss) vengono pulite isolatamente quando specificato.
Usalo --dry-run per vedere esattamente quali tabelle e quante righe sarebbero interessate prima di eseguire una pulizia.
Parametri disponibili per airflow db clean
La tabella seguente descrive i parametri disponibili perairflow db clean.
| Parametro |
Description |
Predefinita |
--clean-before-timestamp |
(Obbligatorio) La data o il timestamp prima del quale i dati vengono eliminati. Se non viene fornito alcun fuso orario, viene utilizzato il fuso orario predefinito di Apache Airflow. Ad esempio: 2025-01-01T00:00:00+00:00 |
Nessuno |
--tables o -t |
Nomi delle tabelle su cui eseguire la manutenzione (separati da virgole). Le opzioni includono: dag_runtask_instance,task_instance_history,log,job,xcom,import_error,,task_reschedule,trigger,dag,dag_version,sla_miss,callback_request,celery_taskmeta,celery_tasksetmeta,asset_event,deadline, revoked_token task_state_store connection_test_request _xcom_archive |
Nessuno |
--batch-size |
Numero massimo di righe da eliminare o archiviare in una singola transazione. Valori più bassi riducono i blocchi di lunga durata ma aumentano il numero di batch. |
Nessuno |
--dry-run |
Esegui una corsa a secco senza eliminare effettivamente i dati. Consigliato per i test iniziali. |
False |
--skip-archive |
Non conservare i record eliminati in una tabella di archivio. Per impostazione predefinita, db clean sposta i record eliminati nelle tabelle di archivio (denominate con una _<table>_archive convenzione, ad esempio_dag_run_archive) anziché eliminarli definitivamente. Ciò fornisce una rete di sicurezza: è possibile ispezionare i dati archiviati, esportarli o rilasciarli in un airflow db export-archived secondo momento. airflow db drop-archived Una volta --skip-archive impostato, i record vengono eliminati definitivamente senza questo passaggio intermedio. |
False |
--dag-ids |
Elimina solo i dati relativi agli ID DAG specificati. |
Nessuno |
--exclude-dag-ids |
Evita di pulire i dati relativi agli ID DAG specificati. |
Nessuno |
-y, --yes |
Salta la richiesta di conferma. Necessario per l'esecuzione non interattiva della CLI tramite Amazon MWAA. |
False |
-v, --verbose |
Rendi l'output di registrazione più dettagliato. |
False |
Per ulteriori informazioni sui parametri disponibili, consultate il riferimento alle variabili CLI e env sul sito Web di Apache Airflow.
Esempi di codice
Gli esempi seguenti mostrano come richiamare airflow db clean tramite l'endpoint CLI di Amazon MWAA. Per ulteriori informazioni sulla creazione di token CLI, consulta. Creazione di un token CLI Apache Airflow
Usando uno 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}")
Esempio di funzionamento a secco (primo passaggio consigliato)
Prima di eseguire una pulizia effettiva, esegui --dry-run per vedere cosa verrebbe eliminato:
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
)