View a markdown version of this page

Pulizia del database Aurora PostgreSQL in un ambiente Amazon MWAA - Amazon Managed Workflows for Apache Airflow

Le traduzioni sono generate tramite traduzione automatica. In caso di conflitto tra il contenuto di una traduzione e la versione originale in Inglese, quest'ultima prevarrà.

Pulizia del database Aurora PostgreSQL in un ambiente Amazon MWAA

Amazon Managed Workflows for Apache Airflow utilizza un database Aurora PostgreSQL come database di metadati Apache Airflow, in cui vengono eseguite le esecuzioni del DAG e vengono archiviate le istanze delle attività. Il seguente codice di esempio cancella periodicamente le voci dal database Aurora PostgreSQL dedicato per il tuo ambiente Amazon MWAA.

Importante

Apache Airflow v3 limita l'accesso diretto al database di metadati dal codice delle attività. I lavoratori non si connettono più al database dei metadati e il DAG o il codice delle attività non possono importare o utilizzare direttamente sessioni o modelli del database Apache Airflow. Questa modifica migliora la sicurezza e la scalabilità. Tuttavia, l'approccio di pulizia del DAG-based database che funziona in Apache Airflow v2 non funziona negli ambienti Apache Airflow v3.

Utilizza invece il comando CLI tramite l'endpoint airflow db clean CLI di Amazon MWAA per eseguire la pulizia del database di metadati.

Nota

Nel tempo, il database di metadati accumula vecchi record, dati XCom e dati di attività obsoleti. Questa crescita utilizza le connessioni al database, rallenta l'ambiente e ritarda le attività. Esegui una pulizia regolare dei metadati per prevenire questi problemi.

Versione

Gli esempi di codice in questa pagina sono specifici per Apache Airflow v2 e v3 supportati su Amazon MWAA. Fai riferimento alle versioni di Apache Airflow supportate. Versioni di Apache Airflow su Flussi di lavoro gestiti da Amazon per Apache Airflow

Prerequisiti

Per utilizzare il codice di esempio in questa pagina, è necessario quanto segue:

Dipendenze

Per utilizzare questo esempio di codice con Apache Airflow v2, non sono richieste dipendenze aggiuntive. Usa https://github.com/aws/amazon-mwaa-docker-images aws-mwaa-docker-images per installare Apache Airflow.

Esempio di codice

Gli esempi seguenti mostrano come pulire il database di metadati nel tuo ambiente Amazon MWAA.

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