如何配置WireMock向Kafka发送回调?是否有Kafka版serveEventListeners?
WireMock 向 Kafka 发送回调的实现方案
WireMock 原生并没有提供针对 Kafka 的内置 serveEventListeners 实现,但可以通过以下两种方式实现向 Kafka 发送回调的需求:
1. 自定义 ServeEventListener 扩展
WireMock 的 ServeEventListener 接口允许监听请求完成事件,触发时执行自定义逻辑(比如发送 Kafka 消息),步骤如下:
实现 Kafka 监听器
编写实现 ServeEventListener 的类,在事件回调中集成 Kafka 生产者逻辑:
import com.github.tomakehurst.wiremock.core.ServeEvent; import com.github.tomakehurst.wiremock.extension.ServeEventListener; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import java.util.Properties; public class KafkaServeEventListener implements ServeEventListener { private final KafkaProducer<String, String> producer; public KafkaServeEventListener() { // 初始化 Kafka 生产者配置 Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); this.producer = new KafkaProducer<>(props); } @Override public void onServeComplete(ServeEvent serveEvent) { // 基于请求/响应构造 Kafka 消息内容 String messageContent = String.format( "Request processed: %s %s | Response status: %d", serveEvent.getRequest().getMethod(), serveEvent.getRequest().getUrl(), serveEvent.getResponse().getStatus() ); // 发送消息到指定 Kafka Topic ProducerRecord<String, String> record = new ProducerRecord<>("wiremock-event-topic", messageContent); producer.send(record); } @Override public void shutdown() { // 关闭生产者资源 producer.close(); } }
注册扩展到 WireMock
- Java 代码启动方式:
import com.github.tomakehurst.wiremock.WireMockServer; import static com.github.tomakehurst.wiremock.core.WireMockConfiguration.wireMockConfig; public class WireMockKafkaSetup { public static void main(String[] args) { WireMockServer server = new WireMockServer(wireMockConfig() .extensions(new KafkaServeEventListener())); server.start(); } }
- 独立版 WireMock 启动方式:
将编译好的扩展 JAR 放入 WireMock 的extensions目录,然后执行启动命令:
java -jar wiremock-jre8-standalone.jar --extensions com.yourpackage.KafkaServeEventListener
2. 借助 Webhook 转发到本地 Kafka 服务(轻量方案)
如果不想编写 Java 扩展,可以利用现有 webhook 监听器,配合本地简易 HTTP 服务转发消息到 Kafka:
配置 Stub 中的 Webhook 监听器
{ "request": { "urlPath": "/something-async", "method": "POST" }, "response": { "status": 200, "body": "{ \"result\": \"SUCCESS BRO\" }", "transformers": ["response-template"] }, "serveEventListeners": [ { "name": "webhook", "parameters": { "method": "POST", "url": "http://localhost:8080/forward-to-kafka", "headers": { "Content-Type": "application/json" }, "body": "{\"requestMethod\":\"{{request.method}}\",\"requestUrl\":\"{{request.url}}\",\"responseStatus\":\"{{response.status}}\"}" } } ] }
本地 HTTP 服务逻辑
编写一个简单的 HTTP 服务(比如用 Python/Node.js),接收 Webhook 请求后调用 Kafka 生产者发送消息,示例 Python 代码:
from flask import Flask, request from kafka import KafkaProducer import json app = Flask(__name__) producer = KafkaProducer( bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8') ) @app.route('/forward-to-kafka', methods=['POST']) def forward_to_kafka(): data = request.get_json() producer.send('wiremock-event-topic', value=data) return {'status': 'ok'}, 200 if __name__ == '__main__': app.run(port=8080)
这个方案无需修改 WireMock 核心逻辑,适合快速验证场景,但依赖额外的本地服务。
内容的提问来源于stack exchange,提问作者chigarow
相关产品推荐
相关产品推荐

