Consideraciones importantes
-
Debe especificar qué tan atrás se realiza la limpieza --clean-before-timestamp para controlar. Utilice una marca de tiempo con formato ISO 8601 (por ejemplo,). 2025-01-01T00:00:00+00:00
-
Le recomendamos que especifique que se limite la limpieza --tables a tablas específicas. Si se omite, el comando limpia todas las tablas compatibles.
-
Comience con un ámbito reducido: utilice primero un --clean-before-timestamp valor anterior (más cercano a la fecha de creación del entorno) y una sola tabla. Esto limita la limpieza solo a los registros más antiguos. Como el comando elimina todo lo que esté antes de la marca de tiempo especificada, el uso de una marca de tiempo más reciente da como resultado un alcance de eliminación mayor. Avanza gradualmente la marca de tiempo a medida que vayas ganando confianza en el proceso.
-
Large-scale la limpieza puede afectar al rendimiento de la base de datos. La eliminación de un gran volumen de registros ejerce presión sobre la base de datos de Aurora PostgreSQL y puede afectar a la capacidad de respuesta del entorno. Utilice el --batch-size parámetro para controlar el tamaño de las transacciones y considere la posibilidad de realizar una limpieza durante los períodos de poco tráfico. Tenga cuidado al ejecutar en entornos de producción.
Comportamiento dependiente de la tabla
Al especificar una tabla con--tables, el comando incluye automáticamente todas las tablas dependientes (secundarias) que tengan relaciones de clave externa con la tabla especificada. Los registros de la tabla secundaria se eliminan primero y, a continuación, los registros de la tabla principal, para cumplir con las restricciones de clave externa. Por ejemplo, especificar --tables dag_run también limpiatask_instance,, task_instance_history xcomtask_state_store, y deadline porque estas tablas hacen referencia dag_run a través de claves externas.
En la tabla siguiente se resumen las cadenas de dependencias.
| Tabla especificada |
Se limpiaron las tablas adicionales (dependientes) |
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 |
Las tablas sin dependientes (comolog,, jobimport_error,sla_miss) se limpian de forma aislada cuando se especifican.
--dry-runUtilícelo para ver exactamente qué tablas y cuántas filas se verían afectadas antes de iniciar una limpieza.
Parámetros disponibles para airflow db clean
En la siguiente tabla se describen los parámetros disponibles paraairflow db clean.
| Parámetro |
Description (Descripción) |
Predeterminado |
--clean-before-timestamp |
(Obligatorio) La fecha o marca de tiempo antes de la cual se purgan los datos. Si no se proporciona ninguna zona horaria, se asume la zona horaria predeterminada de Apache Airflow. Ejemplo: 2025-01-01T00:00:00+00:00 |
Ninguno |
--tables o -t |
Nombres de tablas para realizar el mantenimiento (separados por comas). Las opciones incluyen: dag_run task_instancetask_instance_history,log,job,xcom, import_errortask_reschedule,trigger,dag,dag_version,sla_miss,callback_request, celery_taskmetacelery_tasksetmeta,asset_event,deadline,revoked_token, task_state_store connection_test_request _xcom_archive |
Ninguno |
--batch-size |
Cantidad máxima de filas que se pueden eliminar o archivar en una sola transacción. Los valores más bajos reducen los bloqueos prolongados, pero aumentan la cantidad de lotes. |
Ninguno |
--dry-run |
Realice un simulacro sin eliminar realmente los datos. Se recomienda para las pruebas iniciales. |
False |
--skip-archive |
No conserves los registros purgados en una tabla de archivado. De forma predeterminada, db clean mueve los registros purgados a tablas de archivado (con un nombre según una _<table>_archive convención, por ejemplo_dag_run_archive) en lugar de eliminarlos permanentemente. Esto proporciona una red de seguridad: puede inspeccionar los datos archivados, exportarlos o airflow db export-archived eliminarlos más tarde. airflow db drop-archived Cuando --skip-archive se establece, los registros se eliminan permanentemente sin este paso intermedio. |
False |
--dag-ids |
Solo se limpian los datos relacionados con los ID de DAG dados. |
Ninguno |
--exclude-dag-ids |
Evite borrar los datos relacionados con los ID de DAG dados. |
Ninguno |
-y, --yes |
Omite el mensaje de confirmación. Necesario para la ejecución no interactiva de la CLI mediante Amazon MWAA. |
False |
-v, --verbose |
Haga que la salida del registro sea más detallada. |
False |
Para obtener más información sobre los parámetros disponibles, consulte la referencia a las variables de entorno y CLI en el sitio web de Apache Airflow.
Ejemplos de código
Los siguientes ejemplos muestran cómo invocar a airflow db clean través del punto de enlace de la CLI de Amazon MWAA. Para obtener más información sobre la creación de tokens de CLI, consulte. Creación de un token de la CLI de Apache Airflow
Uso de una secuencia de comandos de 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}")
Ejemplo de ejecución en seco (primer paso recomendado)
Antes de realizar una limpieza real, ejecuta con --dry-run para ver qué se eliminaría:
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
)