Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Задает триггер для потокового запроса. Если он не задан, запрос выполняется как можно быстрее, эквивалентно processingTime='0 seconds'. Одновременно можно задать только один параметр триггера.
Дополнительные сведения см. в разделе "Настройка интервалов триггеров структурированной потоковой передачи".
Синтаксис
trigger(*, processingTime=None, once=None, continuous=None, availableNow=None, realTime=None)
Параметры
| Параметр | Тип | Описание |
|---|---|---|
processingTime |
str, необязательный | Строка интервала обработки (например, '5 seconds', '1 minute'). Периодически выполняет запрос микробатча на основе времени обработки. |
once |
bool, необязательный | Если Trueпри обработке только одного пакета данных завершается запрос. |
continuous |
str, необязательный | Строка интервала времени (например, '5 seconds'). Выполняет непрерывный запрос с заданным интервалом контрольной точки. |
availableNow |
bool, необязательный | Если Trueобработка всех доступных данных в нескольких пакетах завершает запрос. |
realTime |
str, необязательный | Строка длительности пакета (например, '5 seconds'). Выполняет запрос в режиме реального времени с пакетами в указанной длительности. См. режим реального времени в структурированной потоковой передаче. |
Возвраты
DataStreamWriter
Примеры
df = spark.readStream.format("rate").load()
Выполнение триггера каждые 5 секунд:
df.writeStream.trigger(processingTime='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>
Активируйте непрерывное выполнение каждые 5 секунд:
:::примечание о совместимости бессерверных серверов
trigger(continuous=) не поддерживается для бессерверных вычислений Databricks. Для непрерывных конвейеров в бессерверном режиме используйте конвейеры Lakeflow .
:::
df.writeStream.trigger(continuous='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>
Обработка всех доступных данных в нескольких пакетах:
df.writeStream.trigger(availableNow=True)
# <...streaming.readwriter.DataStreamWriter object ...>
Активируйте выполнение в режиме реального времени каждые 5 секунд:
df.writeStream.trigger(realTime='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>