View a markdown version of this page

Aurora PostgreSQL-Datenbankbereinigung in einer Amazon MWAA-Umgebung - Von Amazon verwaltete Workflows für Apache Airflow

Die vorliegende Übersetzung wurde maschinell erstellt. Im Falle eines Konflikts oder eines Widerspruchs zwischen dieser übersetzten Fassung und der englischen Fassung (einschließlich infolge von Verzögerungen bei der Übersetzung) ist die englische Fassung maßgeblich.

Aurora PostgreSQL-Datenbankbereinigung in einer Amazon MWAA-Umgebung

Amazon Managed Workflows for Apache Airflow verwendet eine Aurora PostgreSQL-Datenbank als Apache Airflow-Metadatendatenbank, in der DAG ausgeführt und Task-Instances gespeichert werden. Der folgende Beispielcode löscht regelmäßig Einträge aus der dedizierten Aurora PostgreSQL-Datenbank für Ihre Amazon MWAA-Umgebung.

Wichtig

Apache Airflow v3 schränkt den direkten Zugriff auf Metadaten-Datenbanken über den Taskcode ein. Mitarbeiter stellen keine Verbindung mehr zur Metadatendatenbank her, und DAG oder Taskcode können Apache Airflow-Datenbanksitzungen oder -modelle nicht direkt importieren oder verwenden. Diese Änderung verbessert die Sicherheit und Skalierbarkeit. Der Ansatz zur DAG-based Datenbankbereinigung, der in Apache Airflow v2 funktioniert, funktioniert jedoch nicht in Apache Airflow v3-Umgebungen.

Verwenden Sie stattdessen den airflow db clean CLI-Befehl über den Amazon MWAA-CLI-Endpunkt, um die Metadaten-Datenbank zu bereinigen.

Anmerkung

Im Laufe der Zeit sammelt die Metadaten-Datenbank alte Datensätze, XCom-Daten und veraltete Aufgabendaten an. Dieses Wachstum nutzt Datenbankverbindungen, verlangsamt Ihre Umgebung und verzögert Aufgaben. Führen Sie eine regelmäßige Metadaten-Bereinigung durch, um diese Probleme zu vermeiden.

Version

Die Codebeispiele auf dieser Seite sind spezifisch für Apache Airflow v2 und v3, die auf Amazon MWAA unterstützt werden. Weitere Informationen finden Sie in den unterstützten Apache Airflow-Versionen.

Voraussetzungen

Um den Beispielcode auf dieser Seite verwenden zu können, benötigen Sie Folgendes:

Abhängigkeiten

Um dieses Codebeispiel mit Apache Airflow v2 zu verwenden, sind keine zusätzlichen Abhängigkeiten erforderlich. Verwenden Sie https://github.com/aws/amazon-mwaa-docker-images aws-mwaa-docker-images, um Apache Airflow zu installieren.

Codebeispiel

Die folgenden Beispiele zeigen, wie Sie die Metadaten-Datenbank in Ihrer Amazon MWAA-Umgebung bereinigen.

Apache Airflow v3.0.6 to 3.2.1
Wichtige Überlegungen
  • Sie müssen angeben--clean-before-timestamp, um zu kontrollieren, wie weit die Bereinigung zurückreicht. Verwenden Sie einen Zeitstempel im ISO 8601-Format (z. B.). 2025-01-01T00:00:00+00:00

  • Wir empfehlen, dass Sie angeben, dass die Bereinigung --tables auf bestimmte Tabellen beschränkt werden soll. Wenn er weggelassen wird, bereinigt der Befehl alle unterstützten Tabellen.

  • Beginnen Sie mit einem kleinen Bereich — verwenden Sie einen älteren --clean-before-timestamp Wert (der näher am Erstellungsdatum Ihrer Umgebung liegt) und verwenden Sie zuerst eine einzelne Tabelle. Dadurch wird die Bereinigung nur auf die ältesten Datensätze beschränkt. Da der Befehl alles vor dem angegebenen Zeitstempel löscht, führt die Verwendung eines neueren Zeitstempels zu einem größeren Löschumfang. Verschieben Sie den Zeitstempel schrittweise nach vorne, wenn Sie mehr Vertrauen in den Prozess gewinnen.

  • Large-scale Eine Bereinigung kann sich auf die Datenbankleistung auswirken. Das Löschen einer großen Anzahl von Datensätzen belastet die Aurora PostgreSQL-Datenbank und kann die Reaktionsfähigkeit Ihrer Umgebung beeinträchtigen. Verwenden Sie den --batch-size Parameter, um die Transaktionsgröße zu steuern, und erwägen Sie, die Bereinigung in Zeiten mit geringem Datenverkehr durchzuführen. Seien Sie vorsichtig, wenn Sie es in Produktionsumgebungen ausführen.

Verhalten abhängiger Tabellen

Wenn Sie eine Tabelle mit angeben--tables, schließt der Befehl automatisch alle abhängigen (untergeordneten) Tabellen ein, die Fremdschlüsselbeziehungen zu der angegebenen Tabelle haben. Die Datensätze der untergeordneten Tabelle werden zuerst gelöscht, dann die Datensätze der übergeordneten Tabelle, um Fremdschlüsseleinschränkungen zu erfüllen. Beispielsweise bereinigt die Angabe --tables dag_run auchtask_instance,, und task_instance_history xcomtask_state_store, und deadline weil diese Tabellen dag_run über Fremdschlüssel referenzieren.

In der folgenden Tabelle sind die Abhängigkeitsketten zusammengefasst.

Die angegebene Tabelle Zusätzliche bereinigte Tabellen (Angehörige)
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

Tabellen ohne abhängige Objekte (wielog,,job,sla_miss) werden isoliert bereinigtimport_error, wenn sie angegeben werden.

Verwenden Sie diese Option, --dry-run um genau zu sehen, welche Tabellen und wie viele Zeilen betroffen wären, bevor Sie zu einer Bereinigung übergehen.

Verfügbare Parameter für airflow db clean

In der folgenden Tabelle werden die verfügbaren Parameter für beschriebenairflow db clean.

Parameter Description Standard
--clean-before-timestamp (Erforderlich) Das Datum oder der Zeitstempel, vor dem Daten gelöscht wurden. Wenn keine Zeitzone angegeben wird, wird die Apache Airflow-Standardzeitzone angenommen. Beispiel: 2025-01-01T00:00:00+00:00 Keine
--tables oder -t Tabellennamen, an denen die Wartung durchgeführt werden soll (durch Kommas getrennt). Zu den Optionen gehören: dag_run task_instancetask_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 Keine
--batch-size Maximale Anzahl von Zeilen, die in einer einzigen Transaktion gelöscht oder archiviert werden sollen. Niedrigere Werte reduzieren lang andauernde Sperren, erhöhen aber die Anzahl der Batches. Keine
--dry-run Führen Sie einen Probelauf durch, ohne Daten tatsächlich zu löschen. Für erste Tests empfohlen. Falsch
--skip-archive Bewahren Sie gelöschte Datensätze nicht in einer Archivtabelle auf. db cleanVerschiebt gelöschte Datensätze standardmäßig in Archivtabellen (z. B. nach einer _<table>_archive Konvention benannt_dag_run_archive), anstatt sie dauerhaft zu löschen. Dies bietet ein Sicherheitsnetz — Sie können archivierte Daten überprüfenairflow db export-archived, mit exportieren oder später löschen. airflow db drop-archived Wenn diese Option aktiviert --skip-archive ist, werden Datensätze ohne diesen Zwischenschritt dauerhaft gelöscht. Falsch
--dag-ids Es werden nur Daten bereinigt, die sich auf die angegebenen DAG-IDs beziehen. Keine
--exclude-dag-ids Vermeiden Sie das Bereinigen von Daten, die sich auf die angegebenen DAG-IDs beziehen. Keine
-y, --yes Überspringen Sie die Bestätigungsaufforderung. Erforderlich für die nicht interaktive CLI-Ausführung über Amazon MWAA. Falsch
-v, --verbose Machen Sie die Logging-Ausgabe ausführlicher. Falsch

Weitere Informationen zu den verfügbaren Parametern finden Sie in der Referenz zu CLI- und env-Variablen auf der Apache Airflow-Website.

Codebeispiele

Die folgenden Beispiele zeigen, wie der Aufruf airflow db clean über den Amazon MWAA-CLI-Endpunkt erfolgt. Weitere Informationen zum Erstellen von CLI-Token finden Sie unter. Erstellen eines Apache Airflow CLI-Tokens

Mit einem Python-Skript:

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}")
Beispiel für einen Probelauf (empfohlener erster Schritt)

Bevor Sie eine eigentliche Bereinigung durchführen, führen Sie Folgendes aus, --dry-run um zu sehen, was gelöscht werden würde:

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 )