requestIdExtractor,
String clientId, String groupId, String topic,
boolean autoCommit, int autoCommitIntervalMs) {
Properties props = settings.toProps();
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java
index 7e24ad00a6..2611e9491a 100644
--- a/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java
@@ -1,12 +1,12 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
- *
+ *
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@@ -27,13 +27,12 @@ import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.header.Header;
-import java.nio.ByteBuffer;
-import java.util.Arrays;
import java.util.List;
import java.util.Properties;
import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Future;
-import java.util.function.BiConsumer;
/**
* Created by ashvayka on 24.09.18.
@@ -48,7 +47,7 @@ public class TBKafkaProducerTemplate {
private TbKafkaEnricher enricher = ((value, responseTopic, requestId) -> value);
private final TbKafkaPartitioner partitioner;
- private List partitionInfoList;
+ private ConcurrentMap> partitionInfoMap;
@Getter
private final String defaultTopic;
@@ -78,11 +77,16 @@ public class TBKafkaProducerTemplate {
log.trace("Failed to create topic: {}", e.getMessage(), e);
}
//Maybe this should not be cached, but we don't plan to change size of partitions
- this.partitionInfoList = producer.partitionsFor(defaultTopic);
+ this.partitionInfoMap = new ConcurrentHashMap<>();
+ this.partitionInfoMap.putIfAbsent(defaultTopic, producer.partitionsFor(defaultTopic));
}
- public T enrich(T value, String responseTopic, UUID requestId) {
- return enricher.enrich(value, responseTopic, requestId);
+ T enrich(T value, String responseTopic, UUID requestId) {
+ if (enricher != null) {
+ return enricher.enrich(value, responseTopic, requestId);
+ } else {
+ return value;
+ }
}
public Future send(String key, T value) {
@@ -101,7 +105,7 @@ public class TBKafkaProducerTemplate {
byte[] data = encoder.encode(value);
ProducerRecord record;
Integer partition = getPartition(topic, key, value, data);
- record = new ProducerRecord<>(this.defaultTopic, partition, timestamp, key, data, headers);
+ record = new ProducerRecord<>(topic, partition, timestamp, key, data, headers);
return producer.send(record);
}
@@ -109,7 +113,7 @@ public class TBKafkaProducerTemplate {
if (partitioner == null) {
return null;
} else {
- return partitioner.partition(this.defaultTopic, key, value, data, partitionInfoList);
+ return partitioner.partition(topic, key, value, data, partitionInfoMap.computeIfAbsent(topic, producer::partitionsFor));
}
}
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaHandler.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaHandler.java
new file mode 100644
index 0000000000..66d53c3bde
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaHandler.java
@@ -0,0 +1,12 @@
+package org.thingsboard.server.kafka;
+
+import java.util.function.Consumer;
+
+/**
+ * Created by ashvayka on 05.10.18.
+ */
+public interface TbKafkaHandler {
+
+ void handle(Request request, Consumer onSuccess, Consumer onFailure);
+
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java
index 8a0f5293d9..30b20e709d 100644
--- a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java
@@ -93,13 +93,12 @@ public class TbKafkaRequestTemplate {
ConsumerRecords responses = responseTemplate.poll(Duration.ofMillis(pollInterval));
responses.forEach(response -> {
Header requestIdHeader = response.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER);
- Response decocedResponse = null;
+ Response decodedResponse = null;
UUID requestId = null;
if (requestIdHeader == null) {
try {
- decocedResponse = responseTemplate.decode(response);
- requestId = responseTemplate.extractRequestId(decocedResponse);
-
+ decodedResponse = responseTemplate.decode(response);
+ requestId = responseTemplate.extractRequestId(decodedResponse);
} catch (IOException e) {
log.error("Failed to decode response", e);
}
@@ -107,17 +106,17 @@ public class TbKafkaRequestTemplate {
requestId = bytesToUuid(requestIdHeader.value());
}
if (requestId == null) {
- log.error("[{}] Missing requestId in header and response", response);
+ log.error("[{}] Missing requestId in header and body", response);
} else {
ResponseMetaData expectedResponse = pendingRequests.remove(requestId);
if (expectedResponse == null) {
log.trace("[{}] Invalid or stale request", requestId);
} else {
try {
- if (decocedResponse == null) {
- decocedResponse = responseTemplate.decode(response);
+ if (decodedResponse == null) {
+ decodedResponse = responseTemplate.decode(response);
}
- expectedResponse.future.set(decocedResponse);
+ expectedResponse.future.set(decodedResponse);
} catch (IOException e) {
expectedResponse.future.setException(e);
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java
new file mode 100644
index 0000000000..536c4b93b2
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java
@@ -0,0 +1,173 @@
+/**
+ * Copyright © 2016-2018 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.kafka;
+
+import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.SettableFuture;
+import lombok.Builder;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.kafka.clients.admin.CreateTopicsResult;
+import org.apache.kafka.clients.admin.NewTopic;
+import org.apache.kafka.clients.consumer.ConsumerRecords;
+import org.apache.kafka.common.header.Header;
+import org.apache.kafka.common.header.internals.RecordHeader;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeoutException;
+
+/**
+ * Created by ashvayka on 25.09.18.
+ */
+@Slf4j
+public class TbKafkaResponseTemplate {
+
+ private final TBKafkaConsumerTemplate requestTemplate;
+ private final TBKafkaProducerTemplate responseTemplate;
+ private final TbKafkaHandler handler;
+ private final ConcurrentMap pendingRequests;
+ private final ExecutorService executor;
+ private final long maxPendingRequests;
+
+ private final long pollInterval;
+ private volatile boolean stopped = false;
+
+ @Builder
+ public TbKafkaResponseTemplate(TBKafkaConsumerTemplate requestTemplate,
+ TBKafkaProducerTemplate responseTemplate,
+ TbKafkaHandler handler,
+ long pollInterval,
+ long maxPendingRequests,
+ ExecutorService executor) {
+ this.requestTemplate = requestTemplate;
+ this.responseTemplate = responseTemplate;
+ this.handler = handler;
+ this.pendingRequests = new ConcurrentHashMap<>();
+ this.maxPendingRequests = maxPendingRequests;
+ this.pollInterval = pollInterval;
+ this.executor = executor;
+ }
+
+ public void init() {
+ this.responseTemplate.init();
+ requestTemplate.subscribe();
+ executor.submit(() -> {
+ long nextCleanupMs = 0L;
+ while (!stopped) {
+ ConsumerRecords requests = requestTemplate.poll(Duration.ofMillis(pollInterval));
+ requests.forEach(request -> {
+ Header requestIdHeader = request.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER);
+ if (requestIdHeader == null) {
+ log.error("[{}] Missing requestId in header", request);
+ return;
+ }
+ UUID requestId = bytesToUuid(requestIdHeader.value());
+ if (requestId == null) {
+ log.error("[{}] Missing requestId in header and body", request);
+ return;
+ }
+ Header responseTopicHeader = request.headers().lastHeader(TbKafkaSettings.RESPONSE_TOPIC_HEADER);
+ if (responseTopicHeader == null) {
+ log.error("[{}] Missing response topic in header", request);
+ return;
+ }
+ String responseTopic = bytesToUuid(responseTopicHeader.value());
+ if (requestId == null) {
+ log.error("[{}] Missing requestId in header and body", request);
+ return;
+ }
+
+ Request decodedRequest = null;
+ String responseTopic = null;
+
+ try {
+ if (decodedRequest == null) {
+ decodedRequest = requestTemplate.decode(request);
+ }
+ executor.submit(() -> {
+ handler.handle(decodedRequest, );
+ });
+ } catch (IOException e) {
+ expectedRequest.future.setException(e);
+ }
+
+ });
+ }
+ });
+ }
+
+ public void stop() {
+ stopped = true;
+ }
+
+ public ListenableFuture post(String key, Request request) {
+ if (tickSize > maxPendingRequests) {
+ return Futures.immediateFailedFuture(new RuntimeException("Pending request map is full!"));
+ }
+ UUID requestId = UUID.randomUUID();
+ List headers = new ArrayList<>(2);
+ headers.add(new RecordHeader(TbKafkaSettings.REQUEST_ID_HEADER, uuidToBytes(requestId)));
+ headers.add(new RecordHeader(TbKafkaSettings.RESPONSE_TOPIC_HEADER, stringToBytes(responseTemplate.getTopic())));
+ SettableFuture future = SettableFuture.create();
+ pendingRequests.putIfAbsent(requestId, new ResponseMetaData<>(tickTs + maxRequestTimeout, future));
+ request = requestTemplate.enrich(request, responseTemplate.getTopic(), requestId);
+ requestTemplate.send(key, request, headers);
+ return future;
+ }
+
+ private byte[] uuidToBytes(UUID uuid) {
+ ByteBuffer buf = ByteBuffer.allocate(16);
+ buf.putLong(uuid.getMostSignificantBits());
+ buf.putLong(uuid.getLeastSignificantBits());
+ return buf.array();
+ }
+
+ private static UUID bytesToUuid(byte[] bytes) {
+ ByteBuffer bb = ByteBuffer.wrap(bytes);
+ long firstLong = bb.getLong();
+ long secondLong = bb.getLong();
+ return new UUID(firstLong, secondLong);
+ }
+
+ private byte[] stringToBytes(String string) {
+ return string.getBytes(StandardCharsets.UTF_8);
+ }
+
+ private String bytesToString(byte[] data) {
+ return new String(data, StandardCharsets.UTF_8);
+ }
+
+ private static class ResponseMetaData {
+ private final long expTime;
+ private final SettableFuture future;
+
+ ResponseMetaData(long ts, SettableFuture future) {
+ this.expTime = ts;
+ this.future = future;
+ }
+ }
+
+}
diff --git a/common/transport/pom.xml b/common/transport/pom.xml
index b8a7ab53f0..46efbf4dcc 100644
--- a/common/transport/pom.xml
+++ b/common/transport/pom.xml
@@ -78,6 +78,11 @@
org.springframework
spring-context