View a markdown version of this page

Limpieza de bases de datos de Aurora PostgreSQL en un entorno de Amazon MWAA - Amazon Managed Workflows para Apache Airflow

Las traducciones son generadas a través de traducción automática. En caso de conflicto entre la traducción y la version original de inglés, prevalecerá la version en inglés.

Limpieza de bases de datos de Aurora PostgreSQL en un entorno de Amazon MWAA

Amazon Managed Workflows para Apache Airflow utiliza una base de datos Aurora PostgreSQL como base de datos de metadatos de Apache Airflow, donde se ejecuta el DAG y se almacenan las instancias de tareas. La siguiente muestra de código borra periódicamente las entradas de la base de datos Aurora PostgreSQL dedicada a su entorno Amazon MWAA.

importante

Apache Airflow v3 restringe el acceso directo a la base de datos de metadatos desde el código de tareas. Los trabajadores ya no se conectan a la base de datos de metadatos, y el DAG o el código de tareas no pueden importar ni usar directamente las sesiones o los modelos de la base de datos de Apache Airflow. Este cambio mejora la seguridad y la escalabilidad. Sin embargo, el enfoque DAG-based de limpieza de bases de datos que funciona en Apache Airflow v2 no funciona en los entornos de Apache Airflow v3.

En su lugar, utilice el comando de la airflow db clean CLI a través del punto de enlace de la CLI de Amazon MWAA para limpiar la base de datos de metadatos.

nota

Con el tiempo, la base de datos de metadatos acumula registros antiguos, datos de xCOM y datos de tareas obsoletos. Este crecimiento consume conexiones a bases de datos, ralentiza el entorno y retrasa las tareas. Realice una limpieza regular de los metadatos para evitar estos problemas.

Versión

Los ejemplos de código de esta página son específicos de Apache Airflow v2 y v3 compatibles con Amazon MWAA. Consulte las versiones de Apache Airflow compatibles.

Requisitos previos

Para usar el código de muestra de esta página, necesitará lo siguiente:

Dependencias

Para usar este código de ejemplo con Apache Airflow v2, no se necesitan dependencias adicionales. Use aws-mwaa-docker-images para instalar Apache Airflow.

Código de ejemplo

Los siguientes ejemplos muestran cómo limpiar la base de datos de metadatos de su entorno de Amazon MWAA.

Apache Airflow v3.0.6 to 3.2.1
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"
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 )