View a markdown version of this page

Pembersihan database Aurora PostgreSQL di lingkungan Amazon MWAA - Amazon Managed Workflows for Apache Airflow

Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.

Pembersihan database Aurora PostgreSQL di lingkungan Amazon MWAA

Alur Kerja Terkelola Amazon untuk Apache Airflow menggunakan database Aurora PostgreSQL sebagai database metadata Apache Airflow, tempat DAG berjalan dan instans tugas disimpan. Kode sampel berikut secara berkala menghapus entri dari database Aurora PostgreSQL khusus untuk lingkungan Amazon MWAA Anda.

penting

Apache Airflow v3 membatasi akses database metadata langsung dari kode tugas. Pekerja tidak lagi terhubung ke database metadata, dan DAG atau kode tugas tidak dapat mengimpor atau menggunakan sesi atau model database Apache Airflow secara langsung. Perubahan ini meningkatkan keamanan dan skalabilitas. Namun, pendekatan pembersihan DAG-based database yang berfungsi di Apache Airflow v2 tidak berfungsi di lingkungan Apache Airflow v3.

Sebagai gantinya, gunakan perintah airflow db clean CLI melalui titik akhir Amazon MWAA CLI untuk melakukan pembersihan database metadata.

catatan

Seiring waktu, database metadata mengakumulasi catatan lama, data XCom, dan data tugas basi. Pertumbuhan ini menggunakan koneksi database, memperlambat lingkungan Anda, dan menunda tugas. Jalankan pembersihan metadata reguler untuk mencegah masalah ini.

Versi

Contoh kode di halaman ini khusus untuk Apache Airflow v2 dan v3 yang didukung di Amazon MWAA. Lihat versi Apache Airflow yang didukung.

Prasyarat

Untuk menggunakan kode sampel di halaman ini, Anda memerlukan yang berikut ini:

Dependensi

Untuk menggunakan contoh kode ini dengan Apache Airflow v2, tidak diperlukan dependensi tambahan. Gunakan aws-mwaa-doc ker-images untuk menginstal Apache Airflow.

Contoh kode

Contoh berikut menunjukkan cara membersihkan database metadata di lingkungan Amazon MWAA Anda.

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