如何为WireMock构建简易Kafka模拟服务?
如何为WireMock构建简易Kafka模拟服务?
刚接触Kafka又要搭这种结合WireMock的模拟服务确实有点懵,我之前做过类似的性能测试前置服务,给你梳理个清晰的起步方案,应该能帮你快速上手:
一、先理清核心逻辑
你的模拟服务本质就是Kafka消息消费者 + HTTP请求发送器:
- 从指定Kafka Topic消费消息
- 把消息转换成ServiceA需要的请求格式
- 调用WireMock模拟的ServiceA接口
用Spring Boot来搭会最省心,生态完善,少写很多重复代码。
二、依赖准备(Maven为例)
先在pom.xml里加必要的依赖:
<!-- Spring Kafka 用于消费消息 --> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <!-- Spring Web 自带RestTemplate用来发HTTP请求 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <!-- WireMock 用来模拟ServiceA --> <dependency> <groupId>com.github.tomakehurst</groupId> <artifactId>wiremock-jre8</artifactId> <scope>test</scope> <!-- 如果是生产用的模拟服务就去掉scope --> </dependency>
三、配置Kafka消费者
在application.properties里加基础的Kafka配置:
# Kafka 服务器地址 spring.kafka.bootstrap-servers=localhost:9092 # 消费者组ID spring.kafka.consumer.group-id=mock-service-group # 消息key/value的反序列化器 spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer # 首次启动时从最开始消费消息(方便测试) spring.kafka.consumer.auto-offset-reset=earliest
四、编写Kafka消费者逻辑
写一个监听类,处理收到的Kafka消息,然后调用WireMock:
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import org.springframework.web.client.RestTemplate; @Component public class KafkaMessageConsumer { private final RestTemplate restTemplate; // WireMock的地址,默认是localhost:8080,你可以自己改 private static final String WIREMOCK_SERVICE_A_URL = "http://localhost:8080/service-a/api"; public KafkaMessageConsumer(RestTemplate restTemplate) { this.restTemplate = restTemplate; } @KafkaListener(topics = "your-target-topic", containerFactory = "kafkaListenerContainerFactory") public void consumeMessage(String message) { try { // 这里可以根据实际需求把Kafka消息转换成ServiceA需要的请求体 String serviceARequestBody = convertKafkaMessageToServiceARequest(message); // 调用WireMock模拟的ServiceA接口 String response = restTemplate.postForObject(WIREMOCK_SERVICE_A_URL, serviceARequestBody, String.class); // 打印日志方便调试 System.out.println("Received Kafka message: " + message); System.out.println("Sent request to ServiceA (WireMock), response: " + response); } catch (Exception e) { // 处理异常,比如重试或者记录错误日志 System.err.println("Failed to process Kafka message: " + message + ", error: " + e.getMessage()); } } // 自定义消息转换方法,根据你的实际业务逻辑实现 private String convertKafkaMessageToServiceARequest(String kafkaMessage) { // 示例:直接返回原消息,实际要根据ServiceA的接口格式转换 return kafkaMessage; } }
五、配置WireMock模拟ServiceA
写一个WireMock的初始化类,或者在测试类里启动WireMock并配置Stub:
import com.github.tomakehurst.wiremock.WireMockServer; import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.event.EventListener; import org.springframework.stereotype.Component; import static com.github.tomakehurst.wiremock.client.WireMock.*; @Component public class WireMockServiceASetup { private WireMockServer wireMockServer; @EventListener(ApplicationReadyEvent.class) public void startWireMock() { // 启动WireMock,监听8080端口 wireMockServer = new WireMockServer(8080); wireMockServer.start(); // 配置Stub:当收到POST请求到/service-a/api时,返回指定响应 wireMockServer.stubFor(post(urlEqualTo("/service-a/api")) .willReturn(aResponse() .withStatus(200) .withHeader("Content-Type", "application/json") .withBody("{\"status\":\"success\",\"message\":\"Processed by WireMock\"}"))); } }
六、实用建议
- 本地测试Kafka:如果本地没有Kafka,可以用Testcontainers自动启动Kafka容器,不用手动安装配置,非常方便。
- 消息并发处理:如果性能测试需要高并发,可以调整Kafka消费者的并发数(在@KafkaListener里加
concurrency参数)。 - 异常重试:可以用Spring Kafka的重试机制,比如在配置里加
spring.kafka.listener.retry.max-attempts参数,避免消息丢失。 - 日志优化:用SLF4J代替System.out,方便生产环境排查问题。
先把这个基础版本跑通,然后再根据你的实际业务需求调整消息转换逻辑、WireMock的Stub配置就行,慢慢来,不难的!
备注:内容来源于stack exchange,提问作者kita
相关产品推荐
相关产品推荐

