Flux 객체 외부에서 데이터를 밀어넣을 수 있도록 하는 인터페이스이다.
Flux에 들어갈 데이터소스가 고정된 데이터가 아니라, 그자체도 파이프라인같은 이벤트 리스너나 콜백등일때 사용함.
// 일반적인 Flux - 데이터 소스가 명확할 때
Flux.just(1, 2, 3)
Flux.fromIterable(list)
Flux.interval(Duration.ofSeconds(1))
// 외부 이벤트(콜백, 리스너 등)를 Flux로 변환하고 싶을 때?
// → FluxSink 사용
이벤트 리스너를 Flux로 변환
// 버튼 클릭 같은 외부 이벤트를 Flux로 감싸기
val buttonClicks = Flux.create<String> { sink ->
eventBus.addListener { event ->
sink.next(event.data) // 이벤트 발생할 때마다 주입
}
sink.onDispose {
eventBus.removeListener() // 구독 취소 시 정리
}
}
SSE에서 활용
// 외부에서 데이터를 밀어넣어 SSE로 스트리밍
class NotificationService {
private lateinit var sink: FluxSink<String>
val stream: Flux<String> = Flux.create { sink = it }
// 외부에서 호출해서 데이터 주입
fun send(message: String) {
sink.next(message)
}
}
// 컨트롤러에서
fun stream(request: ServerRequest): Mono<ServerResponse> {
return ServerResponse.ok()
.contentType(MediaType.TEXT_EVENT_STREAM)
.body(notificationService.stream, String::class.java)
}
// 어디서든 호출하면 SSE로 전송됨
notificationService.send("새 알림!")
주요 메서드
Flux.create<String> { sink ->
sink.next("값") // 값 발행
sink.complete() // 완료 신호
sink.error(Exception()) // 에러 신호
sink.onDispose { } // 구독 취소 시 콜백
sink.isCancelled // 구독이 취소됐는지 확인
}
// create - 외부에서 비동기적으로 여러 값을 밀어넣을 때
Flux.create<String> { sink ->
asyncApi.onData { sink.next(it) } // 비동기 콜백
}
// generate - 동기적으로 하나씩 생성할 때
Flux.generate<Int> { sink ->
sink.next(Random.nextInt()) // 한 번에 next 하나만 가능
}