StreamingQueryManager

에 연결된 모든 활성 StreamingQuery 인스턴스를 SparkSession관리합니다. 이 액세스에 사용합니다 spark.streams .

문법

# Access through SparkSession
spark.streams

속성

재산 설명
active 이 SparkSession쿼리와 연결된 모든 활성 스트리밍 쿼리 목록을 반환합니다.

메서드

메서드 설명
get(id) 고유 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()