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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.22 12:00:31