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 kumpulan data
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.
- Pandangan yang 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, data tabel streaming dan status pos pemeriksaan harus tetap ada di Amazon S3. Misalnya, tabel Hive atau AWS Glue-managed (non-Iceberg) mengharuskan database disetel
LocationUrike jalur Amazon S3, sementara tabel Apache Iceberg mengelola sendiri metadata tabelnya.
penting
Jika Anda menggunakan database: bidang di YAML pipeline untuk tabel Hive atau AWS Glue-managed (non-Iceberg), AWS Glue database yang sesuai harus disetel ke jalur LocationUri Amazon S3. Inilah LocationUri yang menempatkan tabel streaming terkelola (dan _spark_metadata lognya) di Amazon S3, yang memungkinkan pemrosesan inkremental cross-run yang persisten untuk tabel tersebut. Tabel gunung es yang menggunakan Katalog Data dengan catalog-impl=GlueCatalog (Opsi 1) tidak memerlukan databaseLocationUri. 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 | AWS Glue Database untuk tabel keluaran. Untuk LocationUri persyaratan, lihatPrasyarat. |
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=...).
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:
- Default (tidak ada bendera run-mode)
-
Menjalankan semua kumpulan data. Tampilan terwujud sepenuhnya dihitung ulang; tabel streaming hanya memproses data baru sejak pos pemeriksaan terakhir.
--refresh <dataset>-
Menyegarkan hanya dataset yang ditentukan. Tabel streaming memproses data baru secara bertahap; tampilan yang terwujud sepenuhnya dihitung ulang.
--full-refresh <dataset>-
Mengatur ulang dan menghitung ulang hanya dataset yang ditentukan. Untuk tabel streaming, ini mengatur ulang pos pemeriksaan dan memproses ulang semua data.
--full-refresh-all-
Mengatur ulang dan menghitung ulang semua kumpulan data.
Menggunakan tabel Iceberg dengan SDP
Apache Iceberg adalah format tabel yang direkomendasikan untuk tabel streaming yang memerlukan pemrosesan inkremental cross-run yang tahan lama, karena tidak bergantung pada _spark_metadata log berbasis file yang digunakan oleh Hive atau tabel yang dikelola. AWS Glue Anda dapat mengonfigurasi Iceberg dengan SDP dalam dua cara, tergantung pada apakah Anda ingin tabel keluaran Anda terdaftar di Katalog Data.
Opsi 1: Gunung Es dengan Katalog Data (disarankan)
Gunakan opsi ini jika Anda ingin tabel Iceberg Anda terdaftar di Katalog Data (dengantable_type=ICEBERG) sehingga dapat dikueri dari mesin lain seperti, Amazon Redshift, dan Amazon EMR. Data tabel dan metadata disimpan di Amazon S3, dan status inkremental cross-run dipertahankan. Tambahkan yang berikut ini ke configuration bagian Andaspark-pipeline.yml:
configuration: spark.sql.catalog.glue_catalog: "org.apache.iceberg.spark.SparkCatalog" spark.sql.catalog.glue_catalog.catalog-impl: "org.apache.iceberg.aws.glue.GlueCatalog" spark.sql.catalog.glue_catalog.io-impl: "org.apache.iceberg.aws.s3.S3FileIO" spark.sql.catalog.glue_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
Di saluran YAML Anda, atur catalog: glue_catalog dan atur database: ke AWS Glue database. Saat Anda membuat AWS Glue pekerjaan, setel --enable-spark-declarative-pipeline ketrue. Jangan mengatur --enable-glue-datacatalog meja gunung es.
catatan
Katalog Iceberg menggunakan lokasinya sendiri warehouse untuk data tabel dan metadata. Bila Anda menggunakan opsi ini, Anda tidak perlu mengatur spark.sql.warehouse.dir atau database LocationUri untuk tabel Iceberg itu sendiri.
Opsi 2: Gunung es dengan katalog berbasis file (Hadoop)
Gunakan opsi ini jika Anda tidak memerlukan pendaftaran Katalog Data. Metadata gunung es berbasis file di Amazon S3, dan tabel tidak terdaftar di Katalog Data. Tambahkan yang berikut ini ke configuration bagian Andaspark-pipeline.yml:
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"
catatan
Perhatikan hal berikut tentang opsi ini:
-
Tabel yang dibuat dengan opsi ini tidak terdaftar di Katalog Data, sehingga mereka tidak dapat dikueri dari mesin kueri seperti.
-
Opsi ini menggunakan lo
warehousekasinya sendiri untuk penyimpanan meja.
Gunakan Opsi 1 jika Anda ingin tabel Anda terdaftar di Katalog Data dan dapat dikueri dari mesin lain. Gunakan Opsi 2 hanya jika Anda tidak memerlukan pendaftaran Katalog Data.
catatan
Parameter --enable-glue-datacatalog pekerjaan menyambungkan metastore Spark Hive ke Katalog Data untuk tabel Hive (non-Iceberg). Untuk tabel Iceberg, catalog-impl=GlueCatalog daftarkan tabel langsung di Katalog Data melalui AWS SDK, jadi Anda tidak menyetel --enable-glue-datacatalog untuk Iceberg. Jangan mengkonfigurasi Iceberg dengan default SparkSessionCatalog (type: hive) bersama dengan --enable-glue-datacatalog dalam upaya untuk mendaftarkan tabel Iceberg di Katalog Data: pada AWS Glue 6.0, kombinasi ini gagal.
Dengan Iceberg 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
Anda juga dapat membaca secara bertahap dari tabel Iceberg sebagai sumber streaming. Dalam arsitektur medali, tabel streaming hilir hanya dapat menggunakan baris baru yang dikomitmen ke tabel Iceberg hulu pada setiap proses. Contoh berikut dibaca secara bertahap dari tabel Iceberg ke bronze tabel silver streaming:
from pyspark import pipelines as dp from pyspark.sql import SparkSession spark = SparkSession.active() dp.create_streaming_table( "silver", comment="Incremental silver layer built from the bronze Iceberg table" ) @dp.append_flow(target="silver") def from_bronze(): # Incremental read from the Iceberg bronze table; each run processes only new rows. return spark.readStream.table("bronze")
Pertimbangan dan batasan
Pertimbangkan hal berikut saat Anda menggunakan SDP:
-
Tampilan yang terwujud selalu dihitung ulang sepenuhnya. Refresh inkremental tidak didukung. Gunakan tabel streaming untuk beban kerja tambahan.
-
Tabel streaming Python API. Gunakan
dp.create_streaming_table()dengan@dp.append_flow(target=...). -
Cross-run pemrosesan tambahan untuk tabel streaming. Tabel streaming mendukung pemrosesan inkremental cross-run hanya jika data dan status pos pemeriksaan tetap ada di Amazon S3. Tabel Hive atau AWS Glue-managed (non-Iceberg) mengharuskan database diset
LocationUriel ke jalur Amazon S3, sementara tabel Apache Iceberg mengelola sendiri metadata tabelnya. -
Database LocationUri diperlukan untuk tabel Hive atau AWS Glue-managed (non-Iceberg). Tabel gunung es yang menggunakan Katalog Data dengan
catalog-impl=GlueCatalogtidak memerlukannya. Lihat perinciannya di Prasyarat. -
Harapan kualitas data. Anotasi kualitas data sebaris tidak didukung dalam kerangka SDP saat ini.
-
Hin
withColumndari fungsi kueri hilir. Ketika kumpulan data hilir (seperti tampilan terwujud) membaca dari kumpulan data pipeline hulu menggunakanspark.table(...)dan diterapkan.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 perkakas 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 dengan mengubah satu
spark.sql(...).write.saveAsTable(...)panggilan 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.