Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Saluran Pipa Deklaratif Spark
Spark Declarative Pipelines (SDP) adalah kerangka kerja deklaratif untuk membangun pipeline data batch dan streaming di 6.0. AWS Glue Dengan SDP, Anda menentukan seperti apa data Anda seharusnya menggunakan SQL atau Python, dan kerangka kerja secara otomatis menentukan rencana eksekusi, menyelesaikan dependensi antara kumpulan data, dan menjalankan cabang independen secara paralel.
SDP menyederhanakan pengembangan pipeline dengan menghilangkan kode boilerplate imperatif untuk membaca, menulis, pendaftaran katalog, dan pemesanan eksekusi. Anda fokus pada transformasi bisnis sementara kerangka kerja menangani infrastruktur pipa.
SDP tersedia dalam AWS Glue versi 6.0 dan yang lebih baru.
Konsep SDP
Pipeline terdiri dari file manifes YAML (spark-pipeline.yml) dan satu atau lebih file transformasi SQL atau Python. SDP secara otomatis:
-
Menyelesaikan dependensi dengan menyimpulkan DAG dari referensi tabel
-
Menentukan urutan eksekusi tanpa orkestrasi manual
-
Menjalankan cabang independen secara paralel untuk throughput maksimum
-
Mengelola status tambahan untuk streaming tabel melalui pos pemeriksaan
-
Mendaftarkan tabel keluaran dalam katalog setelah materialisasi
Jenis dataset
Tiga jenis dataset tersedia di SDP:
- Tabel streaming
-
Hanya memproses data baru sejak proses terakhir. Mempertahankan status di seluruh eksekusi pekerjaan menggunakan pos pemeriksaan. Gunakan tabel streaming untuk penyerapan, aliran peristiwa, data IoT, perubahan pengambilan data, dan penambahan sumber saja.
- Tampilan terwujud
-
Menghitung ulang dataset sepenuhnya pada setiap proses. Output selalu mencerminkan keadaan data sumber saat ini. Gunakan tampilan yang terwujud untuk agregasi, gabungan, analitik ringkasan, dan laporan.
- Tampilan sementara
-
Session-scoped dan tidak dipertahankan atau dikatalogkan. Gunakan tampilan sementara untuk transformasi menengah dan logika pementasan.
penting
Tampilan terwujud selalu melakukan komputasi ulang penuh dalam versi saat ini. Mereka tidak mendukung penyegaran tambahan. Gunakan tabel streaming untuk beban kerja tambahan.
Prasyarat
Untuk menggunakan SDP, Anda memerlukan yang berikut:
-
AWS Glue versi 6.0
-
Lokasi Amazon S3 untuk penyimpanan pipeline (pos pemeriksaan, metadata)
-
Untuk integrasi Katalog Data (opsional): disetel
--enable-glue-datacatalogketrue. Atau, Anda dapat mengonfigurasi pengaturan katalog langsung melalui konfigurasi Spark. -
Untuk penyimpanan tabel persisten: setel
spark.sql.warehouse.dirke jalur Amazon S3, atau aturdatabase:bidang di saluran YAML dan pastikan AWS Glue database telahLocationUridikonfigurasi ke jalur Amazon S3 -
Untuk pemrosesan inkremental cross-run dengan tabel streaming: gunakan tabel Iceberg (tabel Hive-managed streaming tidak mendukung pemrosesan inkremental cross-run)
penting
Jika Anda menggunakan database: bidang di YAML pipeline Anda, AWS Glue database yang sesuai harus diset LocationUri el ke jalur Amazon S3. Database yang dibuat melalui konsol sering kosongLocationUri. Buat atau perbarui database dengan lokasi Amazon S3 eksplisit:
aws glue create-database --database-input '{ "Name":"my_pipeline_db", "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db" }'
Membuat pipa
Untuk membuat pipeline SDP, selesaikan langkah-langkah berikut.
Langkah 1: Buat pipa YAML
Buat file bernama spark-pipeline.yml:
name: my_analytics_pipeline catalog: spark_catalog database: analytics_db storage: s3://my-bucket/pipeline-storage/ libraries: - glob: include: transformations/** configuration: spark.sql.shuffle.partitions: "4"
Tabel berikut menjelaskan bidang YAML pipeline.
| Bidang | Diperlukan | Deskripsi |
|---|---|---|
name |
Ya | Nama untuk pipeline Anda. |
catalog |
Tidak | Katalog untuk digunakan. Default ke spark_catalog. |
database |
Tidak | Basis data target untuk tabel keluaran. Database harus ada dan memiliki satu LocationUri set. |
storage |
Ya | Jalur Amazon S3 untuk pos pemeriksaan pipa dan metadata. |
libraries |
Ya | Pola glob untuk file transformasi untuk disertakan. |
configuration |
Tidak | Properti konfigurasi percikan. |
Langkah 2: Tulis transformasi
Buat file transformasi dalam transformations/ direktori. Anda dapat menggunakan SQL, Python, atau keduanya dalam pipeline yang sama.
Contoh SQL (transformations/silver.sql):
CREATE MATERIALIZED VIEW silver_sales AS SELECT *, UPPER(region) as clean_region FROM bronze_sales WHERE amount > 0; CREATE MATERIALIZED VIEW gold_summary AS SELECT clean_region, COUNT(*) as order_count, SUM(amount) as total_revenue FROM silver_sales GROUP BY clean_region;
Contoh Python (transformations/bronze.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() @dp.materialized_view(comment="Raw sales data from S3") def bronze_sales() -> DataFrame: return spark.read.format("csv").option("header", "true") \ .option("inferSchema", "true") \ .load("s3://source-bucket/raw-data/sales/")
Contoh tabel streaming Python (transformations/events.py):
from pyspark import pipelines as dp from pyspark.sql import DataFrame, SparkSession spark = SparkSession.active() dp.create_streaming_table( "streaming_events", comment="Incremental event ingestion", schema="event_id STRING, event_type STRING, timestamp LONG, payload STRING" ) @dp.append_flow(target="streaming_events") def ingest_events() -> DataFrame: return ( spark.readStream.format("json") .schema("event_id STRING, event_type STRING, timestamp LONG, payload STRING") .load("s3://source-bucket/events/") )
catatan
Untuk streaming tabel di Python, gunakan dp.create_streaming_table() gabungan dengan@dp.append_flow(target=...). @dp.streaming_tableDekorator tidak tersedia.
Langkah 3: Unggah ke Amazon S3
Unggah file pipeline Anda ke Amazon S3 sebagai berikut:
-
.zipFile yang berisispark-pipeline.ymldantransformations/direktori -
Awalan Amazon S3 (direktori) yang berisi struktur yang sama
Langkah 4: Buat dan jalankan AWS Glue pekerjaan
Buat AWS Glue pekerjaan dengan parameter berikut:
-
--enable-spark-declarative-pipeline:true(diperlukan - mengaktifkan mode SDP) -
ScriptLocation: zip definisi pipeline atau awalan Amazon S3 (diperlukan untuk pipeline SDP) -
--enable-glue-datacatalog:true(opsional — mendaftarkan tabel di Katalog Data)
Contoh berikut membuat pekerjaan SDP menggunakan AWS CLI:
aws glue create-job \ --name my-sdp-pipeline \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 2 \ --command '{"Name":"glueetl","ScriptLocation":"s3://my-bucket/pipelines/my_pipeline.zip"}' \ --default-arguments '{ "--enable-spark-declarative-pipeline": "true", "--enable-glue-datacatalog": "true" }'
Menjalankan saluran pipa
Anda menjalankan pipeline SDP menggunakanStartJobRun. Anda dapat mengontrol perilaku eksekusi dengan argumen pekerjaan yang diteruskan pada waktu berjalan.
Mode Jalankan
Berikan argumen berikut StartJobRun untuk mengontrol eksekusi pipeline:
--conf spark.glue.sdp.jobMode-
Mengontrol mode eksekusi:
RUN(default) — Menjalankan pipeline secara normal.VALIDATE- Melakukan dry run yang memeriksa sintaks YAML, resolusi ketergantungan, dan SQL/Python kompilasi tanpa menulis data apa pun.
--conf spark.glue.sdp.runMode-
Mengontrol kumpulan data mana yang diperbarui:
--refresh-
Menjalankan semua kumpulan data. Tampilan yang terwujud sepenuhnya dihitung ulang. Tabel streaming hanya memproses data baru sejak pos pemeriksaan terakhir.
--refresh <dataset_name>-
Hanya menjalankan dataset yang ditentukan. Untuk tabel streaming, ini memproses data baru secara bertahap. Untuk tampilan yang terwujud, ini melakukan penghitungan ulang penuh dari tampilan itu saja.
--full-refresh-
Mengatur ulang dan menghitung ulang semua kumpulan data. Untuk tabel streaming, ini mengatur ulang pos pemeriksaan dan memproses ulang semua data dari awal.
--full-refresh-all-
Menjatuhkan semua tabel dan memproses ulang seluruh pipa dari awal.
Menggunakan tabel Iceberg dengan SDP
Untuk tabel streaming yang memerlukan pemrosesan inkremental cross-run, gunakan Apache Iceberg. Hive-managed tabel streaming menyimpan metadata secara lokal dan tidak bertahan di seluruh proses pekerjaan.
Untuk mengkonfigurasi Iceberg, tambahkan yang berikut ini ke bagian spark-pipeline.yml konfigurasi Anda:
configuration: spark.sql.catalog.spark_catalog: "org.apache.iceberg.spark.SparkSessionCatalog" spark.sql.catalog.spark_catalog.type: "hadoop" spark.sql.catalog.spark_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
Dengan Iceberg yang dikonfigurasi, Anda mendapatkan manfaat berikut:
-
Tabel streaming mempertahankan status pos pemeriksaan di Amazon S3 di seluruh proses pekerjaan
-
Setiap proses membuat file data baru dan snapshot gunung es
-
Proses selanjutnya dilanjutkan dari offset yang dikomit terakhir
-
Riwayat tabel lengkap dipertahankan melalui mekanisme snapshot Iceberg
Pertimbangan dan batasan
Pertimbangkan hal berikut saat Anda menggunakan SDP:
-
Tampilan terwujud selalu dihitung ulang sepenuhnya — Penyegaran tambahan tidak didukung. Gunakan tabel streaming untuk beban kerja tambahan.
-
Tabel streaming Python API - Gunakan
dp.create_streaming_table()dengan@dp.append_flow(target=...).@dp.streaming_tableDekorator tidak tersedia dalam versi saat ini. -
Cross-run pemrosesan tambahan memerlukan Iceberg — Tabel streaming dengan Hive atau katalog ter AWS Glue kelola tidak mendukung pemrosesan tambahan di seluruh proses pekerjaan. Gunakan tabel Iceberg untuk status inkremental persisten.
-
Database LocationUri diperlukan — Jika Anda menentukan YAML
database:dalam pipeline, AWS Glue database harus disetLocationUriel ke jalur Amazon S3. Tanpa itu, pipa gagal. -
Harapan kualitas data — Anotasi kualitas data sebaris tidak didukung dalam kerangka SDP saat ini.
-
Hindari
withColumndalam fungsi kueri hilir — Ketika kumpulan data hilir (seperti tampilan terwujud) membaca dari kumpulan data pipeline hulu menggunakanspark.table(...)dan menerapkan.withColumn(...), SDP mungkin gagal mendeteksi ketergantungan antara kumpulan data pada proses kedua dan berikutnya. Hal ini menyebabkan hilir membaca data basi dari proses sebelumnya (jeda sekali jalan). Untuk menghindari masalah ini, ekspresikan kolom turunan di dalam al.select(...)ih-alih menggunakan.withColumn(...). Hindari juga operasi apa pun yang memaksa resolusi rencana (seperti.schemaatau.collect) di dalam fungsi kueri. -
Tidak ada alat migrasi — Migrasi otomatis dari kerangka kerja pipeline lainnya tidak didukung. Migrasikan tabel secara bertahap — SDP dapat membaca dari tabel katalog yang ada.
-
Penjadwalan — Pekerjaan SDP menggunakan mekanisme penjadwalan yang sama dengan AWS Glue pekerjaan lain (Pem AWS Glue icu, Amazon EventBridge, Apache Airflow).
Migrasi dari skrip imperatif ke SDP
Anda dapat memigrasikan skrip Spark imperatif yang ada ke SDP secara bertahap:
-
Mulailah dengan satu tabel — ubah
spark.sql(...).write.saveAsTable(...)panggilan tunggal menjadi pernyataanCREATE MATERIALIZED VIEWSQL. -
Tambahkan tabel secara bertahap — SDP menangani dependensi campuran. Tabel SDP dapat membaca dari tabel katalog yang ada yang bukan bagian dari pipeline.
-
Jalankan kedua pola secara paralel selama transisi — pekerjaan SDP dan pekerjaan imperatif dapat hidup berdampingan.
SDP dapat mereferensikan tabel apa pun yang dapat diakses melalui SparkSession, termasuk tabel Katalog Data yang ada, tabel eksternal, dan referensi lintas database.