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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:20:44