StreamingQuery

Дескриптор запроса, выполняющегося непрерывно в фоновом режиме по мере поступления новых данных. Все методы являются потокобезопасными.

Синтаксис

# Returned by DataStreamWriter.start() or DataStreamWriter.toTable()
q = df.writeStream.format("console").start()

Свойства

Недвижимость Описание
id Возвращает уникальный идентификатор этого запроса, который сохраняется во время перезапуска из данных контрольной точки.
runId Возвращает уникальный идентификатор этого запроса, который не сохраняется во время перезапуска.
name Возвращает имя запроса, указанное пользователем, или None , если оно не указано.
isActive Возвращает, активен ли этот запрос потоковой передачи.
status Возвращает текущее состояние запроса в качестве дикта.
recentProgress Возвращает массив последних StreamingQueryProgress обновлений для этого запроса.
lastProgress Возвращает последнее StreamingQueryProgress обновление или None нет обновлений.

Методы

Метод Описание
awaitTermination(timeout) Ожидает завершения этого запроса либо по stop() исключению, либо по исключению.
processAllAvailable() Блокирует до тех пор, пока все доступные данные в источнике не будут обработаны и зафиксированы в приемнике. Предназначено для тестирования.
stop() Останавливает этот запрос потоковой передачи.
explain(extended) Выводит планы (логические и физические) в консоль для отладки.
exception() Возвращает значение, StreamingQueryException если запрос завершился с исключением или None.

Примеры

sdf = spark.readStream.format("rate").load()
sq = sdf.writeStream.format('memory').queryName('this_query').start()
sq.isActive
# True
sq.name
# 'this_query'
sq.awaitTermination(5)
# False
sq.stop()
sq.isActive
# False