如何在Kotlin中测试SseEmitter的事件发送功能
Spring Kotlin SSE控制器自动化测试问题解决
问题背景
已实现基于Spring的Kotlin SSE控制器,核心功能为:客户端通过/api/{orderToken}/add建立SSE连接并存储,调用/api/{orderToken}/notification可向对应客户端推送消息。功能手动验证正常,但自动化测试无法正确捕获推送的消息。
控制器代码
import org.springframework.web.bind.annotation.GetMapping import org.springframework.web.bind.annotation.PathVariable import org.springframework.web.bind.annotation.RequestMapping import org.springframework.web.bind.annotation.RestController import org.springframework.web.servlet.mvc.method.annotation.SseEmitter @RestController @RequestMapping("/api/{orderToken}") class TestingController { companion object { private val store = HashMap<String, SseEmitter>() } @GetMapping("/add", produces = ["text/event-stream"]) fun getOrderStatus( @PathVariable orderToken: String, ): SseEmitter { val sseEmitter = SseEmitter() store[orderToken] = sseEmitter return sseEmitter } @GetMapping("/notification", produces = ["text/event-stream"]) fun simulateKafkaEvent( @PathVariable orderToken: String, ) { val sseEmitter = store[orderToken] sseEmitter!!.send("order token: $orderToken") } }
原有测试代码问题
原有测试错误尝试消费SseEmitter对象而非SSE推送的消息,同时存在端点路径错误:
@BeforeEach fun setup() { webTestClient = WebTestClient.bindToController(TestingController()).build() } @Test fun `test 2`() { val addEndpoint = webTestClient .get() .uri("/api/1234/add") .accept(MediaType.TEXT_EVENT_STREAM) .exchange() .expectStatus().isOk .expectHeader().contentTypeCompatibleWith(MediaType.TEXT_EVENT_STREAM) .returnResult(SseEmitter::class.java) .responseBody // 错误:SSE响应是消息流,不是SseEmitter对象 val notification = webTestClient.get() .uri("/api/1234/kafka") // 错误:控制器端点为/notification,非/kafka .accept() .accept(MediaType.TEXT_EVENT_STREAM) .exchange() .returnResult(String::class.java) .responseBody StepVerifier.create(addEndpoint).consumeNextWith { println("notification: " + s)}.thenCancel().verify() }
错误分析
- 类型认知错误:
SseEmitter是服务端用于发送消息的内部对象,不会作为响应体返回给客户端;WebTestClient处理SSE时,返回的是推送的消息流(Flux<ServerSentEvent<T>>或直接Flux<T>)。 - 路径错误:触发推送的端点路径与控制器定义不匹配,导致无法触发消息推送。
- 流程顺序错误:未确保SSE连接已建立并订阅就触发推送,可能导致消息丢失。
正确测试代码
import org.junit.jupiter.api.BeforeEach import org.junit.jupiter.api.Test import org.springframework.http.MediaType import org.springframework.test.web.reactive.server.WebTestClient import reactor.test.StepVerifier class TestingControllerTest { private lateinit var webTestClient: WebTestClient @BeforeEach fun setup() { webTestClient = WebTestClient.bindToController(TestingController()).build() } @Test fun `test sse notification delivery`() { val orderToken = "1234" // 1. 建立SSE连接,获取推送消息流 val sseMessageFlux = webTestClient.get() .uri("/api/{orderToken}/add", orderToken) .accept(MediaType.TEXT_EVENT_STREAM) .exchange() .expectStatus().isOk .expectHeader().contentTypeCompatibleWith(MediaType.TEXT_EVENT_STREAM) .returnResult(String::class.java) .responseBody // 2. 触发消息推送 webTestClient.get() .uri("/api/{orderToken}/notification", orderToken) .exchange() .expectStatus().isOk // 3. 验证消息是否推送成功 StepVerifier.create(sseMessageFlux) .consumeNextWith { receivedMessage -> assert(receivedMessage == "order token: $orderToken") println("Received notification: $receivedMessage") } .thenCancel() // 手动取消连接,避免测试阻塞 .verify() } }
关键修正说明
- 响应类型修正:使用
returnResult(String::class.java)直接获取推送的消息流(Flux<String>),对应SSE推送的内容。 - 端点路径修正:将触发推送的路径改为
/notification,与控制器定义一致。 - 流程优化:先建立SSE连接并拿到消息流,再触发推送,确保消息能被捕获。
- 断言验证:添加断言确保收到的消息符合预期,提升测试可靠性。
内容的提问来源于stack exchange,提问作者Boris
相关产品推荐
相关产品推荐

