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