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 yazıcı kullanılarak işlenecek akış sorgusunun çıkışını ayarlar. İşleme mantığı, giriş olarak bir satır alan bir işlev olarak veya ve isteğe bağlı process(row) ve open(partition_id, epoch_id) yöntemlerine sahip close(error) bir nesne olarak belirtilebilir.
Sözdizimi
foreach(f)
Parametreler
| Parametre | Türü | Açıklama |
|---|---|---|
f |
çağrılabilir veya nesne | Giriş olarak Satır alan bir işlev veya yöntemi ve isteğe bağlı process(row) ve open yöntemleri olan bir close nesne. |
İadeler
DataStreamWriter
Notlar
Sağlanan nesne serileştirilebilir olmalıdır. Veri yazmak için herhangi bir başlatma (örneğin, bir bağlantı açma) içinde yapılmalıdır open(), oluşturma zamanında değil.
Örnekler
import time
df = spark.readStream.format("rate").load()
İşlev kullanarak her satırı işleyin:
def print_row(row):
print(row)
q = df.writeStream.foreach(print_row).start()
time.sleep(3)
q.stop()
, ve open yöntemlerine sahip processclosebir nesnesi kullanarak her satırı işleyin:
class RowPrinter:
def open(self, partition_id, epoch_id):
print("Opened %d, %d" % (partition_id, epoch_id))
return True
def process(self, row):
print(row)
def close(self, error):
print("Closed with error: %s" % str(error))
q = df.writeStream.foreach(RowPrinter()).start()
time.sleep(3)
q.stop()