Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Sağlanan işlev kullanılarak işlenecek akış sorgusunun çıkışını ayarlar. Yalnızca mikro toplu yürütme modunda (tetikleyici sürekli olmadığında) desteklenir. Her mikro toplu işlemde, sağlanan işlev çıktı satırları DataFrame ve toplu iş tanımlayıcısı olarak çağrılır. Toplu iş kimliği, çıktıyı yinelenenleri kaldırıp dış sistemlere işlem yoluyla yazmak için kullanılabilir.
Sözdizimi
foreachBatch(func)
Parametreler
| Parametre | Türü | Açıklama |
|---|---|---|
func |
Callable | Giriş olarak DataFrame ve toplu iş kimliği (int) alan bir işlev. |
İadeler
DataStreamWriter
Notlar
Spark Connect modunda, sağlanan işlevin dışında tanımlanan değişkenlere erişimi yoktur.
Örnekler
import time
df = spark.readStream.format("rate").load()
def func(batch_df, batch_id):
batch_df.collect()
q = df.writeStream.foreachBatch(func).start()
time.sleep(3)
q.stop()