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"
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
)