Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Декоратор @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 должен правильно обрабатывать полное обновление для подчиненных источников данных. Дополнительные сведения см. в разделе "Полное обновление". |