foreach_batch_sink

Декоратор @dp.foreach_batch_sink() определяет приемник ForEachBatch, который обрабатывает поток как ряд микропакетов, которые обрабатываются в Python с пользовательской логикой. Вы ссылаетесь на приемник в виде targetпотока добавления для записи преобразованных данных. Концептуальные рекомендации, рекомендации и примеры см. в разделе Use ForEachBatch для записи в произвольные приемники данных в конвейерах.

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.

Parameters

Parameter Описание
name Optional. Уникальное имя для идентификации приемника в конвейере. По умолчанию используется имя UDF, если оно не указано.
batch_handler Это определяемая пользователем функция (UDF), вызываемая для каждого микропакета.
df Кадр данных Spark, содержащий данные для текущего микропакета.
batch_id Целочисленный идентификатор микробатча. Spark увеличивает этот идентификатор для каждого интервала триггера.
batch_id 0 представляет собой начало потока или начало полного обновления. Код foreach_batch_sink должен правильно обрабатывать полное обновление для подчиненных источников данных. Дополнительные сведения см. в разделе "Полное обновление".