StreamingQueryManager

Управляет всеми активными StreamingQuery экземплярами, связанными с объектом SparkSession. Используйте spark.streams для доступа к этому.

Синтаксис

# Access through SparkSession
spark.streams

Свойства

Недвижимость Описание
active Возвращает список всех активных потоковых запросов, связанных с этим SparkSession.

Методы

Метод Описание
get(id) Возвращает активный запрос по уникальному идентификатору.
awaitAnyTermination(timeout) Ожидает завершения любого активного запроса или до истечения срока ожидания.
resetTerminated() Забывает прошлые завершенные запросы, чтобы awaitAnyTermination() его можно было использовать еще раз, чтобы ждать новых завершения.
addListener(listener) Регистрирует обратные StreamingQueryListener вызовы событий жизненного цикла.
removeListener(listener) Отменяет StreamingQueryListenerрегистрацию .

Примеры

sdf = spark.readStream.format("rate").load()
sq = sdf.writeStream.format('memory').queryName('this_query').start()
sqm = spark.streams
[q.name for q in sqm.active]
# ['this_query']
sqm.awaitAnyTermination(5)
# True
sq.stop()
sqm.resetTerminated()