|
|
|
@ -92,6 +92,9 @@ public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTe |
|
|
|
long nextCleanupMs = 0L; |
|
|
|
while (!stopped) { |
|
|
|
ConsumerRecords<String, byte[]> responses = responseTemplate.poll(Duration.ofMillis(pollInterval)); |
|
|
|
if (responses.count() > 0) { |
|
|
|
log.trace("Polling responses completed, consumer records count [{}]", responses.count()); |
|
|
|
} |
|
|
|
responses.forEach(response -> { |
|
|
|
Header requestIdHeader = response.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER); |
|
|
|
Response decodedResponse = null; |
|
|
|
@ -109,6 +112,7 @@ public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTe |
|
|
|
if (requestId == null) { |
|
|
|
log.error("[{}] Missing requestId in header and body", response); |
|
|
|
} else { |
|
|
|
log.trace("[{}] Response received", requestId); |
|
|
|
ResponseMetaData<Response> expectedResponse = pendingRequests.remove(requestId); |
|
|
|
if (expectedResponse == null) { |
|
|
|
log.trace("[{}] Invalid or stale request", requestId); |
|
|
|
@ -132,6 +136,7 @@ public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTe |
|
|
|
if (kv.getValue().expTime < tickTs) { |
|
|
|
ResponseMetaData<Response> staleRequest = pendingRequests.remove(kv.getKey()); |
|
|
|
if (staleRequest != null) { |
|
|
|
log.trace("[{}] Request timeout detected, expTime [{}], tickTs [{}]", kv.getKey(), staleRequest.expTime, tickTs); |
|
|
|
staleRequest.future.setException(new TimeoutException()); |
|
|
|
} |
|
|
|
} |
|
|
|
@ -158,8 +163,10 @@ public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTe |
|
|
|
headers.add(new RecordHeader(TbKafkaSettings.REQUEST_ID_HEADER, uuidToBytes(requestId))); |
|
|
|
headers.add(new RecordHeader(TbKafkaSettings.RESPONSE_TOPIC_HEADER, stringToBytes(responseTemplate.getTopic()))); |
|
|
|
SettableFuture<Response> future = SettableFuture.create(); |
|
|
|
pendingRequests.putIfAbsent(requestId, new ResponseMetaData<>(tickTs + maxRequestTimeout, future)); |
|
|
|
ResponseMetaData<Response> responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future); |
|
|
|
pendingRequests.putIfAbsent(requestId, responseMetaData); |
|
|
|
request = requestTemplate.enrich(request, responseTemplate.getTopic(), requestId); |
|
|
|
log.trace("[{}] Sending request, key [{}], expTime [{}]", requestId, key, responseMetaData.expTime); |
|
|
|
requestTemplate.send(key, request, headers, null); |
|
|
|
return future; |
|
|
|
} |
|
|
|
|