30 changed files with 328 additions and 169 deletions
@ -0,0 +1,102 @@ |
|||||
|
/** |
||||
|
* 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.service.queue; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.dao.queue.MsgQueue; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.atomic.AtomicLong; |
||||
|
|
||||
|
@Service |
||||
|
@Slf4j |
||||
|
public class DefaultMsgQueueService implements MsgQueueService { |
||||
|
|
||||
|
@Value("${rule.queue.max_size}") |
||||
|
private long queueMaxSize; |
||||
|
|
||||
|
@Value("${rule.queue.cleanup_period}") |
||||
|
private long queueCleanUpPeriod; |
||||
|
|
||||
|
@Autowired |
||||
|
private MsgQueue msgQueue; |
||||
|
|
||||
|
private ScheduledExecutorService cleanupExecutor; |
||||
|
|
||||
|
private Map<TenantId, AtomicLong> pendingCountPerTenant = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
if (queueCleanUpPeriod > 0) { |
||||
|
cleanupExecutor = Executors.newSingleThreadScheduledExecutor(); |
||||
|
cleanupExecutor.scheduleAtFixedRate(() -> cleanup(), |
||||
|
queueCleanUpPeriod, queueCleanUpPeriod, TimeUnit.SECONDS); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stop() { |
||||
|
if (cleanupExecutor != null) { |
||||
|
cleanupExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { |
||||
|
AtomicLong pendingMsgCount = pendingCountPerTenant.computeIfAbsent(tenantId, key -> new AtomicLong()); |
||||
|
if (pendingMsgCount.incrementAndGet() < queueMaxSize) { |
||||
|
return msgQueue.put(tenantId, msg, nodeId, clusterPartition); |
||||
|
} else { |
||||
|
pendingMsgCount.decrementAndGet(); |
||||
|
return Futures.immediateFailedFuture(new RuntimeException("Message queue is full!")); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { |
||||
|
ListenableFuture<Void> result = msgQueue.ack(tenantId, msg, nodeId, clusterPartition); |
||||
|
AtomicLong pendingMsgCount = pendingCountPerTenant.computeIfAbsent(tenantId, key -> new AtomicLong()); |
||||
|
pendingMsgCount.decrementAndGet(); |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public Iterable<TbMsg> findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { |
||||
|
return msgQueue.findUnprocessed(tenantId, nodeId, clusterPartition); |
||||
|
} |
||||
|
|
||||
|
private void cleanup() { |
||||
|
pendingCountPerTenant.forEach((tenantId, pendingMsgCount) -> { |
||||
|
pendingMsgCount.set(0); |
||||
|
msgQueue.cleanUp(tenantId); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,32 @@ |
|||||
|
/** |
||||
|
* 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.service.queue; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
public interface MsgQueueService { |
||||
|
|
||||
|
ListenableFuture<Void> put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); |
||||
|
|
||||
|
ListenableFuture<Void> ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition); |
||||
|
|
||||
|
Iterable<TbMsg> findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,118 @@ |
|||||
|
/** |
||||
|
* 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.dao.queue.memory; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.dao.queue.MsgQueue; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import java.util.*; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.concurrent.Executors; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 27.04.18. |
||||
|
*/ |
||||
|
@Component |
||||
|
@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) |
||||
|
@Slf4j |
||||
|
public class InMemoryMsgQueue implements MsgQueue { |
||||
|
|
||||
|
private ListeningExecutorService queueExecutor; |
||||
|
private Map<TenantId, Map<InMemoryMsgKey, Map<UUID, TbMsg>>> data = new HashMap<>(); |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
// Should be always single threaded due to absence of locks.
|
||||
|
queueExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stop() { |
||||
|
if (queueExecutor != null) { |
||||
|
queueExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { |
||||
|
return queueExecutor.submit(() -> { |
||||
|
data.computeIfAbsent(tenantId, key -> new HashMap<>()). |
||||
|
computeIfAbsent(new InMemoryMsgKey(nodeId, clusterPartition), key -> new HashMap<>()).put(msg.getId(), msg); |
||||
|
return null; |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { |
||||
|
return queueExecutor.submit(() -> { |
||||
|
Map<InMemoryMsgKey, Map<UUID, TbMsg>> tenantMap = data.get(tenantId); |
||||
|
if (tenantMap != null) { |
||||
|
InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); |
||||
|
Map<UUID, TbMsg> map = tenantMap.get(key); |
||||
|
if (map != null) { |
||||
|
map.remove(msg.getId()); |
||||
|
if (map.isEmpty()) { |
||||
|
tenantMap.remove(key); |
||||
|
} |
||||
|
} |
||||
|
if (tenantMap.isEmpty()) { |
||||
|
data.remove(tenantId); |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public Iterable<TbMsg> findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { |
||||
|
ListenableFuture<List<TbMsg>> list = queueExecutor.submit(() -> { |
||||
|
Map<InMemoryMsgKey, Map<UUID, TbMsg>> tenantMap = data.get(tenantId); |
||||
|
if (tenantMap != null) { |
||||
|
InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); |
||||
|
Map<UUID, TbMsg> map = tenantMap.get(key); |
||||
|
if (map != null) { |
||||
|
return new ArrayList<>(map.values()); |
||||
|
} else { |
||||
|
return Collections.emptyList(); |
||||
|
} |
||||
|
} else { |
||||
|
return Collections.emptyList(); |
||||
|
} |
||||
|
}); |
||||
|
try { |
||||
|
return list.get(); |
||||
|
} catch (InterruptedException | ExecutionException e) { |
||||
|
throw new RuntimeException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Void> cleanUp(TenantId tenantId) { |
||||
|
return queueExecutor.submit(() -> { |
||||
|
data.remove(tenantId); |
||||
|
return null; |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -1,114 +0,0 @@ |
|||||
/** |
|
||||
* 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.dao.sql.queue; |
|
||||
|
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import com.google.common.util.concurrent.ListeningExecutorService; |
|
||||
import com.google.common.util.concurrent.MoreExecutors; |
|
||||
import lombok.Getter; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
import org.thingsboard.server.common.msg.TbMsg; |
|
||||
import org.thingsboard.server.dao.queue.MsgQueue; |
|
||||
import org.thingsboard.server.dao.util.SqlDao; |
|
||||
|
|
||||
import javax.annotation.PostConstruct; |
|
||||
import javax.annotation.PreDestroy; |
|
||||
import java.util.*; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.atomic.AtomicLong; |
|
||||
|
|
||||
/** |
|
||||
* Created by ashvayka on 27.04.18. |
|
||||
*/ |
|
||||
@Component |
|
||||
//@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true)
|
|
||||
@Slf4j |
|
||||
@SqlDao |
|
||||
public class InMemoryMsgQueue implements MsgQueue { |
|
||||
|
|
||||
@Value("${rule.queue.max_size}") |
|
||||
@Getter |
|
||||
private long maxSize; |
|
||||
|
|
||||
private ListeningExecutorService queueExecutor; |
|
||||
private AtomicLong pendingMsgCount = new AtomicLong(); |
|
||||
private Map<InMemoryMsgKey, Map<UUID, TbMsg>> data = new HashMap<>(); |
|
||||
|
|
||||
@PostConstruct |
|
||||
public void init() { |
|
||||
// Should be always single threaded due to absence of locks.
|
|
||||
queueExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); |
|
||||
} |
|
||||
|
|
||||
@PreDestroy |
|
||||
public void stop() { |
|
||||
if (queueExecutor == null) { |
|
||||
queueExecutor.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> put(TbMsg msg, UUID nodeId, long clusterPartition) { |
|
||||
if (pendingMsgCount.incrementAndGet() < maxSize) { |
|
||||
return queueExecutor.submit(() -> { |
|
||||
data.computeIfAbsent(new InMemoryMsgKey(nodeId, clusterPartition), key -> new HashMap<>()).put(msg.getId(), msg); |
|
||||
return null; |
|
||||
}); |
|
||||
} else { |
|
||||
pendingMsgCount.decrementAndGet(); |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("Message queue is full!")); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> ack(TbMsg msg, UUID nodeId, long clusterPartition) { |
|
||||
return queueExecutor.submit(() -> { |
|
||||
InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); |
|
||||
Map<UUID, TbMsg> map = data.get(key); |
|
||||
if (map != null) { |
|
||||
if (map.remove(msg.getId()) != null) { |
|
||||
pendingMsgCount.decrementAndGet(); |
|
||||
} |
|
||||
if (map.isEmpty()) { |
|
||||
data.remove(key); |
|
||||
} |
|
||||
} |
|
||||
return null; |
|
||||
}); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Iterable<TbMsg> findUnprocessed(UUID nodeId, long clusterPartition) { |
|
||||
ListenableFuture<List<TbMsg>> list = queueExecutor.submit(() -> { |
|
||||
InMemoryMsgKey key = new InMemoryMsgKey(nodeId, clusterPartition); |
|
||||
Map<UUID, TbMsg> map = data.get(key); |
|
||||
if (map != null) { |
|
||||
return new ArrayList<>(map.values()); |
|
||||
} else { |
|
||||
return Collections.emptyList(); |
|
||||
} |
|
||||
}); |
|
||||
try { |
|
||||
return list.get(); |
|
||||
} catch (InterruptedException | ExecutionException e) { |
|
||||
throw new RuntimeException(e); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
Loading…
Reference in new issue