You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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()
}

错误分析

  1. 类型认知错误:SseEmitter是服务端用于发送消息的内部对象,不会作为响应体返回给客户端;WebTestClient处理SSE时,返回的是推送的消息流(Flux<ServerSentEvent<T>>或直接Flux<T>)。
  2. 路径错误:触发推送的端点路径与控制器定义不匹配,导致无法触发消息推送。
  3. 流程顺序错误:未确保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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 16:33:18