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       // 구독이 취소됐는지 확인
}

Flux.create vs Flux.generate

// create - 외부에서 비동기적으로 여러 값을 밀어넣을 때
Flux.create<String> { sink ->
    asyncApi.onData { sink.next(it) }  // 비동기 콜백
}

// generate - 동기적으로 하나씩 생성할 때
Flux.generate<Int> { sink ->
    sink.next(Random.nextInt())  // 한 번에 next 하나만 가능
}