Browse Source

Transport

pull/1166/head
Andrew Shvayka 8 years ago
parent
commit
b190a4bd40
  1. 14
      application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java
  2. 7
      application/src/main/java/org/thingsboard/server/service/transport/TransportApiService.java
  3. 50
      common/queue/src/main/java/org/thingsboard/server/kafka/AbstractTbKafkaTemplate.java
  4. 29
      common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java
  5. 98
      common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java

14
application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java

@ -0,0 +1,14 @@
package org.thingsboard.server.service.transport;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service;
/**
* Created by ashvayka on 05.10.18.
*/
@Slf4j
@Service
@ConditionalOnProperty(prefix = "quota.rule.tenant", value = "enabled", havingValue = "true", matchIfMissing = false)
public class RemoteTransportApiService implements TransportApiService {
}

7
application/src/main/java/org/thingsboard/server/service/transport/TransportApiService.java

@ -0,0 +1,7 @@
package org.thingsboard.server.service.transport;
/**
* Created by ashvayka on 05.10.18.
*/
public interface TransportApiService {
}

50
common/queue/src/main/java/org/thingsboard/server/kafka/AbstractTbKafkaTemplate.java

@ -0,0 +1,50 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
* <p>
* 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
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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 lombok.extern.slf4j.Slf4j;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.UUID;
/**
* Created by ashvayka on 25.09.18.
*/
@Slf4j
public abstract class AbstractTbKafkaTemplate {
protected byte[] uuidToBytes(UUID uuid) {
ByteBuffer buf = ByteBuffer.allocate(16);
buf.putLong(uuid.getMostSignificantBits());
buf.putLong(uuid.getLeastSignificantBits());
return buf.array();
}
protected static UUID bytesToUuid(byte[] bytes) {
ByteBuffer bb = ByteBuffer.wrap(bytes);
long firstLong = bb.getLong();
long secondLong = bb.getLong();
return new UUID(firstLong, secondLong);
}
protected byte[] stringToBytes(String string) {
return string.getBytes(StandardCharsets.UTF_8);
}
protected String bytesToString(byte[] data) {
return new String(data, StandardCharsets.UTF_8);
}
}

29
common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java

@ -23,24 +23,25 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.CreateTopicsResult;
import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.header.internals.RecordHeader; import org.apache.kafka.common.header.internals.RecordHeader;
import java.io.IOException; import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.time.Duration; import java.time.Duration;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.*; 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. * Created by ashvayka on 25.09.18.
*/ */
@Slf4j @Slf4j
public class TbKafkaRequestTemplate<Request, Response> { public class TbKafkaRequestTemplate<Request, Response> extends AbstractTbKafkaTemplate {
private final TBKafkaProducerTemplate<Request> requestTemplate; private final TBKafkaProducerTemplate<Request> requestTemplate;
private final TBKafkaConsumerTemplate<Response> responseTemplate; private final TBKafkaConsumerTemplate<Response> responseTemplate;
@ -163,24 +164,6 @@ public class TbKafkaRequestTemplate<Request, Response> {
return future; 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 static class ResponseMetaData<T> { private static class ResponseMetaData<T> {
private final long expTime; private final long expTime;
private final SettableFuture<T> future; private final SettableFuture<T> future;

98
common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaResponseTemplate.java

@ -15,52 +15,45 @@
*/ */
package org.thingsboard.server.kafka; 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.Builder;
import lombok.extern.slf4j.Slf4j; 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.clients.consumer.ConsumerRecords;
import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.header.internals.RecordHeader; import org.apache.kafka.common.header.internals.RecordHeader;
import java.io.IOException; import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.time.Duration; import java.time.Duration;
import java.util.ArrayList; import java.util.Collections;
import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.TimeoutException;
/** /**
* Created by ashvayka on 25.09.18. * Created by ashvayka on 25.09.18.
*/ */
@Slf4j @Slf4j
public class TbKafkaResponseTemplate<Request, Response> { public class TbKafkaResponseTemplate<Request, Response> extends AbstractTbKafkaTemplate {
private final TBKafkaConsumerTemplate<Request> requestTemplate; private final TBKafkaConsumerTemplate<Request> requestTemplate;
private final TBKafkaProducerTemplate<Response> responseTemplate; private final TBKafkaProducerTemplate<Response> responseTemplate;
private final TbKafkaHandler<Request, Response> handler; private final TbKafkaHandler<Request, Response> handler;
private final ConcurrentMap<UUID, String> pendingRequests; private final ConcurrentMap<UUID, String> pendingRequests;
private final ExecutorService executor; private final ExecutorService executor;
private final long maxPendingRequests; private final int maxPendingRequests;
private final long pollInterval; private final long pollInterval;
private volatile boolean stopped = false; private volatile boolean stopped = false;
//TODO:
private final AtomicInteger pendingRequestCount = new AtomicInteger();
@Builder @Builder
public TbKafkaResponseTemplate(TBKafkaConsumerTemplate<Request> requestTemplate, public TbKafkaResponseTemplate(TBKafkaConsumerTemplate<Request> requestTemplate,
TBKafkaProducerTemplate<Response> responseTemplate, TBKafkaProducerTemplate<Response> responseTemplate,
TbKafkaHandler<Request, Response> handler, TbKafkaHandler<Request, Response> handler,
long pollInterval, long pollInterval,
long maxPendingRequests, int maxPendingRequests,
ExecutorService executor) { ExecutorService executor) {
this.requestTemplate = requestTemplate; this.requestTemplate = requestTemplate;
this.responseTemplate = responseTemplate; this.responseTemplate = responseTemplate;
@ -75,8 +68,11 @@ public class TbKafkaResponseTemplate<Request, Response> {
this.responseTemplate.init(); this.responseTemplate.init();
requestTemplate.subscribe(); requestTemplate.subscribe();
executor.submit(() -> { executor.submit(() -> {
long nextCleanupMs = 0L;
while (!stopped) { while (!stopped) {
if(pendingRequestCount.get() > maxPendingRequests){
}
//TODO: we need to protect from reading too much requests.
ConsumerRecords<String, byte[]> requests = requestTemplate.poll(Duration.ofMillis(pollInterval)); ConsumerRecords<String, byte[]> requests = requestTemplate.poll(Duration.ofMillis(pollInterval));
requests.forEach(request -> { requests.forEach(request -> {
Header requestIdHeader = request.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER); Header requestIdHeader = request.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER);
@ -94,26 +90,15 @@ public class TbKafkaResponseTemplate<Request, Response> {
log.error("[{}] Missing response topic in header", request); log.error("[{}] Missing response topic in header", request);
return; return;
} }
String responseTopic = bytesToUuid(responseTopicHeader.value()); String responseTopic = bytesToString(responseTopicHeader.value());
if (requestId == null) {
log.error("[{}] Missing requestId in header and body", request);
return;
}
Request decodedRequest = null;
String responseTopic = null;
try { try {
if (decodedRequest == null) { Request decodedRequest = requestTemplate.decode(request);
decodedRequest = requestTemplate.decode(request); executor.submit(() -> handler.handle(decodedRequest,
} response -> reply(requestId, responseTopic, response),
executor.submit(() -> { e -> log.error("[{}] Failed to process the request: {}", requestId, request, e)));
handler.handle(decodedRequest, ); } catch (Throwable e) {
}); log.error("[{}] Failed to process the request: {}", requestId, request, e);
} catch (IOException e) {
expectedRequest.future.setException(e);
} }
}); });
} }
}); });
@ -123,51 +108,8 @@ public class TbKafkaResponseTemplate<Request, Response> {
stopped = true; stopped = true;
} }
public ListenableFuture<Response> post(String key, Request request) { private void reply(UUID requestId, String topic, Response response) {
if (tickSize > maxPendingRequests) { responseTemplate.send(topic, response, Collections.singletonList(new RecordHeader(TbKafkaSettings.REQUEST_ID_HEADER, uuidToBytes(requestId))));
return Futures.immediateFailedFuture(new RuntimeException("Pending request map is full!"));
}
UUID requestId = UUID.randomUUID();
List<Header> 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<Response> 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<T> {
private final long expTime;
private final SettableFuture<T> future;
ResponseMetaData(long ts, SettableFuture<T> future) {
this.expTime = ts;
this.future = future;
}
} }
} }

Loading…
Cancel
Save