Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Dekorator @dp.foreach_batch_sink() mendefinisikan sink ForEachBatch, yang memproses aliran sebagai serangkaian batch mikro yang Anda tangani di Python dengan logika kustom. Anda mereferensikan sink sebagai target dalam alur tambahan untuk menulis data yang diubah. Untuk panduan konseptual, pertimbangan, dan contoh, lihat Menggunakan ForEachBatch untuk menulis ke sink data arbitrer dalam alur.
Syntax
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="<name>")
def batch_handler(df, batch_id):
"""
Required:
- `df`: a Spark DataFrame representing the rows of this micro-batch.
- `batch_id`: unique integer ID for each micro-batch in the query.
"""
# Your custom write or transformation logic here
# Example:
# df.write.format("some-target-system").save("...")
#
# To access the sparkSession inside the batch handler, use df.sparkSession.
Parameter-parameternya
| Parameter | Deskripsi |
|---|---|
| nama | Optional. Nama unik untuk mengidentifikasi sink dalam alur. Diatur ke nama UDF secara default, ketika tidak disertakan. |
| batch_handler | Ini adalah fungsi yang ditentukan pengguna (UDF) yang dipanggil untuk setiap mikro-batch. |
| Df | Spark DataFrame yang berisi data untuk mikro-batch saat ini. |
| batch_id | ID bilangan bulat dari mikro-batch. Spark menaikkan ID ini untuk setiap interval pemicu. Dari batch_id0 mewakili awal aliran, atau awal refresh penuh. Kode foreach_batch_sink harus menangani pembaruan penuh untuk sumber data hilir dengan tepat. Untuk informasi selengkapnya, lihat Refresh penuh. |