Pertimbangan penting
-
Anda harus menentukan --clean-before-timestamp untuk mengontrol seberapa jauh jarak pembersihan mencapai. Gunakan stempel waktu berformat ISO 8601 (misalnya,). 2025-01-01T00:00:00+00:00
-
Sebaiknya tentukan --tables untuk membatasi pembersihan pada tabel tertentu. Jika dihilangkan, perintah membersihkan semua tabel yang didukung.
-
Mulailah dengan cakupan kecil — gunakan --clean-before-timestamp nilai yang lebih lama (lebih dekat dengan tanggal pembuatan lingkungan Anda) dan satu tabel terlebih dahulu. Ini membatasi pembersihan hanya pada catatan tertua. Karena perintah menghapus semuanya sebelum stempel waktu yang ditentukan, menggunakan stempel waktu yang lebih baru menghasilkan cakupan penghapusan yang lebih besar. Secara bertahap pindahkan stempel waktu ke depan saat Anda mendapatkan kepercayaan diri dalam prosesnya.
-
Large-scale pembersihan dapat memengaruhi kinerja database. Menghapus volume rekaman yang tinggi memberi tekanan pada database Aurora PostgreSQL dan dapat memengaruhi respons lingkungan Anda. Gunakan --batch-size parameter untuk mengontrol ukuran transaksi, dan pertimbangkan untuk menjalankan pembersihan selama periode lalu lintas rendah. Berhati-hatilah saat berjalan di lingkungan produksi.
Perilaku tabel dependen
Ketika Anda menentukan tabel dengan--tables, perintah secara otomatis menyertakan setiap tabel dependen (anak) yang memiliki hubungan kunci asing dengan tabel yang ditentukan. Catatan tabel anak dihapus terlebih dahulu, kemudian catatan tabel induk, untuk memenuhi batasan kunci asing. Misalnya, menentukan --tables dag_run juga membersihkantask_instance,,task_instance_history, xcomtask_state_store, dan deadline karena tabel ini merujuk dag_run melalui kunci asing.
Tabel berikut merangkum rantai ketergantungan.
| Tabel ditentukan |
Tabel tambahan dibersihkan (tanggungan) |
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 |
Tabel tanpa tanggungan (sepertilog,job,import_error,sla_miss) dibersihkan secara terpisah ketika ditentukan.
Gunakan --dry-run untuk melihat dengan tepat tabel mana dan berapa banyak baris yang akan terpengaruh sebelum melakukan pembersihan.
Parameter yang tersedia untuk airflow db clean
Tabel berikut menjelaskan parameter yang tersedia untukairflow db clean.
| Parameter |
Deskripsi |
Default |
--clean-before-timestamp |
(Diperlukan) Tanggal atau stempel waktu sebelum data dibersihkan. Jika tidak ada zona waktu yang disediakan, zona waktu default Apache Airflow diasumsikan. Contoh: 2025-01-01T00:00:00+00:00 |
Tidak ada |
--tables atau -t |
Nama tabel untuk melakukan pemeliharaan (dipisahkan koma). Pilihan meliputi:dag_run,task_instance,task_instance_history,log,job,xcom,import_error, task_rescheduletrigger,dag,dag_version,sla_miss,callback_request,celery_taskmeta,celery_tasksetmeta,asset_event,deadline,revoked_token,task_state_store,connection_test_request, _xcom_archive |
Tidak ada |
--batch-size |
Jumlah maksimum baris yang akan dihapus atau diarsipkan dalam satu transaksi. Nilai yang lebih rendah mengurangi kunci yang berjalan lama tetapi meningkatkan jumlah batch. |
Tidak ada |
--dry-run |
Lakukan dry run tanpa benar-benar menghapus data. Direkomendasikan untuk pengujian awal. |
False |
--skip-archive |
Jangan menyimpan catatan yang dibersihkan dalam tabel arsip. Secara default, db clean memindahkan catatan yang dibersihkan ke dalam tabel arsip (diberi nama dengan _<table>_archive konvensi, misalnya_dag_run_archive) alih-alih menghapusnya secara permanen. Ini menyediakan jaring pengaman — Anda dapat memeriksa data yang diarsipkan, mengekspornyaairflow db export-archived, atau melepaskannya nantiairflow db drop-archived. Ketika --skip-archive disetel, catatan dihapus secara permanen tanpa langkah perantara ini. |
False |
--dag-ids |
Hanya data pembersihan yang terkait dengan ID DAG yang diberikan. |
Tidak ada |
--exclude-dag-ids |
Hindari membersihkan data yang terkait dengan ID DAG yang diberikan. |
Tidak ada |
-y, --yes |
Lewati prompt konfirmasi. Diperlukan untuk eksekusi CLI non-interaktif melalui Amazon MWAA. |
False |
-v, --verbose |
Jadikan keluaran logging lebih bertele-tele. |
False |
Untuk informasi selengkapnya tentang parameter yang tersedia, lihat referensi variabel CLI dan env di situs web Apache Airflow.
Sampel Kode
Contoh berikut menunjukkan cara memanggil airflow db clean melalui titik akhir CLI Amazon MWAA. Untuk informasi selengkapnya tentang membuat token CLI, lihatMembuat token CLI Apache Airflow.
Menggunakan skrip 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}")
Contoh dry run (langkah pertama yang disarankan)
Sebelum melakukan pembersihan yang sebenarnya, jalankan dengan --dry-run untuk melihat apa yang akan dihapus:
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
)