Senior
How would you implement a web service that synchronously sends a message to Kafka and waits for a response from another topic?
sobes.tech AI
Answer from AI
To implement a web service that synchronously sends a message to Kafka and waits for a response from another topic, you can use the following approach:
-
Send a message with a unique correlation ID. This will link the request and response.
-
Subscribe to the response topic with filtering by correlation ID.
-
Use a waiting mechanism with a timeout (e.g., CompletableFuture or CountDownLatch) to block the thread until a response is received or the timeout expires.
Example in Java using Kafka Consumer and Producer:
public String sendAndReceive(String requestTopic, String responseTopic, String message, String correlationId) throws Exception {
// Send message with correlationId in headers
ProducerRecord<String, String> record = new ProducerRecord<>(requestTopic, message);
record.headers().add("correlationId", correlationId.getBytes(StandardCharsets.UTF_8));
producer.send(record);
// Create an object to wait for the response
CompletableFuture<String> futureResponse = new CompletableFuture<>();
// In the consumer, upon receiving a message from responseTopic, check correlationId
consumer.subscribe(Collections.singletonList(responseTopic));
new Thread(() -> {
while (!futureResponse.isDone()) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> rec : records) {
String respCorrelationId = new String(rec.headers().lastHeader("correlationId").value(), StandardCharsets.UTF_8);
if (correlationId.equals(respCorrelationId)) {
futureResponse.complete(rec.value());
return;
}
}
}
}).start();
// Wait for the response with a timeout
return futureResponse.get(10, TimeUnit.SECONDS);
}
This approach provides synchronous behavior over Kafka's asynchronous system, allowing the service to wait for a specific response linked to the sent message.