如何在Vert.x中使用Apache Kafka?含服务端、客户端及微服务通信
在Vert.x中用Apache Kafka实现微服务通信:从基础到请求-响应模式
我懂你的感受——之前用Vert.x Web Client做微服务调用,简单直接,切换到Kafka这种消息驱动的模式,确实得转个弯。毕竟Kafka天生是异步、基于事件的,不像HTTP那样是同步请求响应。下面我就一步步帮你理清怎么在Vert.x里用Kafka,以及怎么实现类似“方法调用”的交互。
第一步:Vert.x集成Kafka的基础配置
首先,你需要引入Vert.x的Kafka客户端依赖(以Maven为例):
<dependency> <groupId>io.vertx</groupId> <artifactId>vertx-kafka-client</artifactId> <version>你的Vert.x版本号</version> </dependency>
1.1 创建Vert.x Kafka生产者
和原生Kafka生产者逻辑一致,但Vert.x版本是异步非阻塞的,更贴合Vert.x的生态:
// 生产者配置 Map<String, String> producerConfig = new HashMap<>(); producerConfig.put("bootstrap.servers", "localhost:9092"); producerConfig.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); producerConfig.put("value.serializer", "io.vertx.kafka.client.serialization.JsonObjectSerializer"); // 创建Vert.x Kafka生产者实例 KafkaProducer<String, JsonObject> producer = KafkaProducer.create(vertx, producerConfig); // 发送示例消息 JsonObject requestMsg = new JsonObject() .put("requestId", UUID.randomUUID().toString()) .put("method", "getUserById") .put("params", new JsonObject().put("userId", 123)); producer.send(new ProducerRecord<>("user-service-requests", requestMsg), ar -> { if (ar.succeeded()) { RecordMetadata metadata = ar.result(); System.out.printf("消息发送成功:主题=%s,分区=%d,偏移量=%d%n", metadata.topic(), metadata.partition(), metadata.offset()); } else { System.err.println("消息发送失败:" + ar.cause().getMessage()); } });
1.2 创建Vert.x Kafka消费者
同样采用异步回调风格,消息到达时自动触发处理逻辑:
// 消费者配置 Map<String, String> consumerConfig = new HashMap<>(); consumerConfig.put("bootstrap.servers", "localhost:9092"); consumerConfig.put("group.id", "user-service-consumer-group"); consumerConfig.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); consumerConfig.put("value.deserializer", "io.vertx.kafka.client.serialization.JsonObjectDeserializer"); consumerConfig.put("auto.offset.reset", "earliest"); // 创建Vert.x Kafka消费者实例 KafkaConsumer<String, JsonObject> consumer = KafkaConsumer.create(vertx, consumerConfig); // 订阅主题并处理消息 consumer.subscribe(Collections.singletonList("user-service-requests"), ar -> { if (ar.succeeded()) { System.out.println("成功订阅请求主题"); consumer.handler(record -> { JsonObject request = record.value(); System.out.println("收到请求:" + request.encodePrettily()); // 调用业务逻辑处理请求 handleServiceRequest(request, producer); }); } else { System.err.println("订阅主题失败:" + ar.cause().getMessage()); } });
第二步:实现类似“方法调用”的请求-响应模式
Kafka本身是单向消息传递,但我们可以通过请求ID+响应主题的方式模拟同步方法调用的效果,核心逻辑是:
- 客户端发送请求时生成唯一
requestId,并指定响应主题 - 服务端处理完请求后,携带同一个
requestId将结果发送到响应主题 - 客户端监听响应主题,通过
requestId匹配对应的请求结果
2.1 客户端(调用方)实现
客户端需要同时扮演生产者(发请求)和消费者(收响应)的角色:
// 初始化响应消费者 Map<String, String> responseConsumerConfig = new HashMap<>(); responseConsumerConfig.put("bootstrap.servers", "localhost:9092"); responseConsumerConfig.put("group.id", "order-service-response-group"); responseConsumerConfig.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); responseConsumerConfig.put("value.deserializer", "io.vertx.kafka.client.serialization.JsonObjectDeserializer"); responseConsumerConfig.put("auto.offset.reset", "earliest"); KafkaConsumer<String, JsonObject> responseConsumer = KafkaConsumer.create(vertx, responseConsumerConfig); responseConsumer.subscribe(Collections.singletonList("user-service-responses"), ar -> { if (ar.succeeded()) { responseConsumer.handler(record -> { JsonObject response = record.value(); String requestId = response.getString("requestId"); // 从待处理请求集合中取出对应的回调 CompletableFuture<JsonObject> pendingRequest = pendingRequests.remove(requestId); if (pendingRequest != null) { if (response.getBoolean("success")) { pendingRequest.complete(response.getJsonObject("result")); } else { pendingRequest.completeExceptionally(new RuntimeException(response.getString("error"))); } } }); } }); // 封装类似方法调用的异步函数 private final Map<String, CompletableFuture<JsonObject>> pendingRequests = new ConcurrentHashMap<>(); public CompletableFuture<JsonObject> callUserService(String method, JsonObject params) { CompletableFuture<JsonObject> future = new CompletableFuture<>(); String requestId = UUID.randomUUID().toString(); JsonObject requestMsg = new JsonObject() .put("requestId", requestId) .put("method", method) .put("params", params) .put("responseTopic", "user-service-responses"); // 指定响应主题 // 存储待处理请求 pendingRequests.put(requestId, future); // 发送请求 producer.send(new ProducerRecord<>("user-service-requests", requestMsg), ar -> { if (ar.failed()) { pendingRequests.remove(requestId); future.completeExceptionally(ar.cause()); } }); // 设置超时机制,避免无限等待 vertx.setTimer(5000, tid -> { if (pendingRequests.containsKey(requestId)) { pendingRequests.remove(requestId); future.completeExceptionally(new TimeoutException("请求超时")); } }); return future; } // 调用示例 callUserService("getUserById", new JsonObject().put("userId", 123)) .thenAccept(user -> System.out.println("获取到用户:" + user.encodePrettily())) .exceptionally(err -> { System.err.println("调用失败:" + err.getMessage()); return null; });
2.2 服务端(提供方)实现
服务端收到请求后处理业务逻辑,再将结果发送回指定的响应主题:
private void handleServiceRequest(JsonObject request, KafkaProducer<String, JsonObject> producer) { String method = request.getString("method"); JsonObject params = request.getJsonObject("params"); String requestId = request.getString("requestId"); String responseTopic = request.getString("responseTopic"); JsonObject response = new JsonObject().put("requestId", requestId); try { // 根据方法名调用对应的业务逻辑 JsonObject result = switch (method) { case "getUserById" -> getUserById(params.getInteger("userId")); case "createUser" -> createUser(params); default -> throw new IllegalArgumentException("不支持的方法:" + method); }; response.put("success", true).put("result", result); } catch (Exception e) { response.put("success", false).put("error", e.getMessage()); } // 发送响应到指定主题 producer.send(new ProducerRecord<>(responseTopic, response), ar -> { if (ar.failed()) { System.err.println("发送响应失败:" + ar.cause().getMessage()); } }); } // 模拟业务方法 private JsonObject getUserById(int userId) { return new JsonObject() .put("id", userId) .put("name", "John Doe") .put("email", "john@example.com"); }
第三步:一些最佳实践
- 主题命名规范:建议用
{服务名}-requests和{服务名}-responses的格式,清晰区分请求与响应主题 - 消息序列化:推荐用JSON(易读)或Protobuf(高性能),Vert.x Kafka客户端已内置JsonObject序列化器,也可自定义
- 错误处理:处理消息发送/消费失败的场景,可结合Kafka重试配置或Vert.x重试逻辑
- 超时机制:务必为请求设置超时,避免客户端长期阻塞
- 消费者组唯一性:每个服务的消费者组ID要唯一,防止重复消费
- 幂等性保障:Kafka可能出现重复消息,服务端需保证业务方法的幂等性(比如用requestId去重)
内容的提问来源于stack exchange,提问作者Isaías Solorio Ayala
相关产品推荐
相关产品推荐

