196 changed files with 3895 additions and 2473 deletions
@ -0,0 +1,17 @@ |
|||||
|
-- |
||||
|
-- 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. |
||||
|
-- |
||||
|
|
||||
|
ALTER TABLE component_descriptor ADD UNIQUE (clazz); |
||||
@ -0,0 +1,39 @@ |
|||||
|
/** |
||||
|
* 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.actors.device; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.SessionType; |
||||
|
|
||||
|
/** |
||||
|
* @author Andrew Shvayka |
||||
|
*/ |
||||
|
@Data |
||||
|
class SessionInfoMetaData { |
||||
|
private final SessionInfo sessionInfo; |
||||
|
private long lastActivityTime; |
||||
|
private boolean subscribedToAttributes; |
||||
|
private boolean subscribedToRPC; |
||||
|
|
||||
|
SessionInfoMetaData(SessionInfo sessionInfo) { |
||||
|
this(sessionInfo, System.currentTimeMillis()); |
||||
|
} |
||||
|
|
||||
|
SessionInfoMetaData(SessionInfo sessionInfo, long lastActivityTime) { |
||||
|
this.sessionInfo = sessionInfo; |
||||
|
this.lastActivityTime = lastActivityTime; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,39 @@ |
|||||
|
/** |
||||
|
* 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.actors.device; |
||||
|
|
||||
|
import org.thingsboard.server.common.msg.MsgType; |
||||
|
import org.thingsboard.server.common.msg.TbActorMsg; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 29.10.18. |
||||
|
*/ |
||||
|
public class SessionTimeoutCheckMsg implements TbActorMsg { |
||||
|
|
||||
|
private static final SessionTimeoutCheckMsg INSTANCE = new SessionTimeoutCheckMsg(); |
||||
|
|
||||
|
private SessionTimeoutCheckMsg() { |
||||
|
} |
||||
|
|
||||
|
public static SessionTimeoutCheckMsg instance() { |
||||
|
return INSTANCE; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public MsgType getMsgType() { |
||||
|
return MsgType.SESSION_TIMEOUT_MSG; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,50 @@ |
|||||
|
/** |
||||
|
* 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.session; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.cache.annotation.CachePut; |
||||
|
import org.springframework.cache.annotation.Cacheable; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Collections; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.CacheConstants.SESSIONS_CACHE; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 29.10.18. |
||||
|
*/ |
||||
|
@Service |
||||
|
@Slf4j |
||||
|
public class DefaultDeviceSessionCacheService implements DeviceSessionCacheService { |
||||
|
|
||||
|
@Override |
||||
|
@Cacheable(cacheNames = SESSIONS_CACHE, key = "#deviceId") |
||||
|
public DeviceSessionsCacheEntry get(DeviceId deviceId) { |
||||
|
log.debug("[{}] Fetching session data from cache", deviceId); |
||||
|
return DeviceSessionsCacheEntry.newBuilder().addAllSessions(Collections.emptyList()).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
@CachePut(cacheNames = SESSIONS_CACHE, key = "#deviceId") |
||||
|
public DeviceSessionsCacheEntry put(DeviceId deviceId, DeviceSessionsCacheEntry sessions) { |
||||
|
log.debug("[{}] Pushing session data from cache: {}", deviceId, sessions); |
||||
|
return sessions; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* 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.session; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 29.10.18. |
||||
|
*/ |
||||
|
public interface DeviceSessionCacheService { |
||||
|
|
||||
|
DeviceSessionsCacheEntry get(DeviceId deviceId); |
||||
|
|
||||
|
DeviceSessionsCacheEntry put(DeviceId deviceId, DeviceSessionsCacheEntry sessions); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,79 @@ |
|||||
|
/** |
||||
|
* 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.nosql; |
||||
|
|
||||
|
import com.datastax.driver.core.ResultSet; |
||||
|
import com.datastax.driver.core.ResultSetFuture; |
||||
|
import com.google.common.util.concurrent.SettableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.scheduling.annotation.Scheduled; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.dao.util.AbstractBufferedRateExecutor; |
||||
|
import org.thingsboard.server.dao.util.AsyncTaskContext; |
||||
|
import org.thingsboard.server.dao.util.NoSqlAnyDao; |
||||
|
|
||||
|
import javax.annotation.PreDestroy; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
@Component |
||||
|
@Slf4j |
||||
|
@NoSqlAnyDao |
||||
|
public class CassandraBufferedRateExecutor extends AbstractBufferedRateExecutor<CassandraStatementTask, ResultSetFuture, ResultSet> { |
||||
|
|
||||
|
public CassandraBufferedRateExecutor( |
||||
|
@Value("${cassandra.query.buffer_size}") int queueLimit, |
||||
|
@Value("${cassandra.query.concurrent_limit}") int concurrencyLimit, |
||||
|
@Value("${cassandra.query.permit_max_wait_time}") long maxWaitTime, |
||||
|
@Value("${cassandra.query.dispatcher_threads:2}") int dispatcherThreads, |
||||
|
@Value("${cassandra.query.callback_threads:2}") int callbackThreads, |
||||
|
@Value("${cassandra.query.poll_ms:50}") long pollMs) { |
||||
|
super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs); |
||||
|
} |
||||
|
|
||||
|
@Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") |
||||
|
public void printStats() { |
||||
|
log.info("Permits queueSize [{}] totalAdded [{}] totalLaunched [{}] totalReleased [{}] totalFailed [{}] totalExpired [{}] totalRejected [{}] currBuffer [{}] ", |
||||
|
getQueueSize(), |
||||
|
totalAdded.getAndSet(0), totalLaunched.getAndSet(0), totalReleased.getAndSet(0), |
||||
|
totalFailed.getAndSet(0), totalExpired.getAndSet(0), totalRejected.getAndSet(0), |
||||
|
concurrencyLevel.get()); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stop() { |
||||
|
super.stop(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected SettableFuture<ResultSet> create() { |
||||
|
return SettableFuture.create(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ResultSetFuture wrap(CassandraStatementTask task, SettableFuture<ResultSet> future) { |
||||
|
return new TbResultSetFuture(future); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ResultSetFuture execute(AsyncTaskContext<CassandraStatementTask, ResultSet> taskCtx) { |
||||
|
CassandraStatementTask task = taskCtx.getTask(); |
||||
|
return task.getSession().executeAsync(task.getStatement()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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.dao.nosql; |
||||
|
|
||||
|
import com.datastax.driver.core.Session; |
||||
|
import com.datastax.driver.core.Statement; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.dao.util.AsyncTask; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
@Data |
||||
|
public class CassandraStatementTask implements AsyncTask { |
||||
|
|
||||
|
private final Session session; |
||||
|
private final Statement statement; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,94 @@ |
|||||
|
/** |
||||
|
* 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.nosql; |
||||
|
|
||||
|
import com.datastax.driver.core.ResultSet; |
||||
|
import com.datastax.driver.core.ResultSetFuture; |
||||
|
import com.google.common.util.concurrent.SettableFuture; |
||||
|
|
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.concurrent.Executor; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.TimeoutException; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
public class TbResultSetFuture implements ResultSetFuture { |
||||
|
|
||||
|
private final SettableFuture<ResultSet> mainFuture; |
||||
|
|
||||
|
public TbResultSetFuture(SettableFuture<ResultSet> mainFuture) { |
||||
|
this.mainFuture = mainFuture; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ResultSet getUninterruptibly() { |
||||
|
return getSafe(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ResultSet getUninterruptibly(long timeout, TimeUnit unit) throws TimeoutException { |
||||
|
return getSafe(timeout, unit); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean cancel(boolean mayInterruptIfRunning) { |
||||
|
return mainFuture.cancel(mayInterruptIfRunning); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean isCancelled() { |
||||
|
return mainFuture.isCancelled(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean isDone() { |
||||
|
return mainFuture.isDone(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ResultSet get() throws InterruptedException, ExecutionException { |
||||
|
return mainFuture.get(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ResultSet get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { |
||||
|
return mainFuture.get(timeout, unit); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void addListener(Runnable listener, Executor executor) { |
||||
|
mainFuture.addListener(listener, executor); |
||||
|
} |
||||
|
|
||||
|
private ResultSet getSafe() { |
||||
|
try { |
||||
|
return mainFuture.get(); |
||||
|
} catch (InterruptedException | ExecutionException e) { |
||||
|
throw new IllegalStateException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private ResultSet getSafe(long timeout, TimeUnit unit) throws TimeoutException { |
||||
|
try { |
||||
|
return mainFuture.get(timeout, unit); |
||||
|
} catch (InterruptedException | ExecutionException e) { |
||||
|
throw new IllegalStateException(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,175 @@ |
|||||
|
/** |
||||
|
* 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.util; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FutureCallback; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.SettableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
|
||||
|
import javax.annotation.Nullable; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.BlockingQueue; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.LinkedBlockingDeque; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.TimeoutException; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
@Slf4j |
||||
|
public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extends ListenableFuture<V>, V> implements BufferedRateExecutor<T, F> { |
||||
|
|
||||
|
private final long maxWaitTime; |
||||
|
private final long pollMs; |
||||
|
private final BlockingQueue<AsyncTaskContext<T, V>> queue; |
||||
|
private final ExecutorService dispatcherExecutor; |
||||
|
private final ExecutorService callbackExecutor; |
||||
|
private final ScheduledExecutorService timeoutExecutor; |
||||
|
private final int concurrencyLimit; |
||||
|
|
||||
|
protected final AtomicInteger concurrencyLevel = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalAdded = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalLaunched = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalReleased = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalFailed = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalExpired = new AtomicInteger(); |
||||
|
protected final AtomicInteger totalRejected = new AtomicInteger(); |
||||
|
|
||||
|
public AbstractBufferedRateExecutor(int queueLimit, int concurrencyLimit, long maxWaitTime, int dispatcherThreads, int callbackThreads, long pollMs) { |
||||
|
this.maxWaitTime = maxWaitTime; |
||||
|
this.pollMs = pollMs; |
||||
|
this.concurrencyLimit = concurrencyLimit; |
||||
|
this.queue = new LinkedBlockingDeque<>(queueLimit); |
||||
|
this.dispatcherExecutor = Executors.newFixedThreadPool(dispatcherThreads); |
||||
|
this.callbackExecutor = Executors.newFixedThreadPool(callbackThreads); |
||||
|
this.timeoutExecutor = Executors.newSingleThreadScheduledExecutor(); |
||||
|
for (int i = 0; i < dispatcherThreads; i++) { |
||||
|
dispatcherExecutor.submit(this::dispatch); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public F submit(T task) { |
||||
|
SettableFuture<V> settableFuture = create(); |
||||
|
F result = wrap(task, settableFuture); |
||||
|
try { |
||||
|
totalAdded.incrementAndGet(); |
||||
|
queue.add(new AsyncTaskContext<>(UUID.randomUUID(), task, settableFuture, System.currentTimeMillis())); |
||||
|
} catch (IllegalStateException e) { |
||||
|
totalRejected.incrementAndGet(); |
||||
|
settableFuture.setException(e); |
||||
|
} |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
public void stop() { |
||||
|
if (dispatcherExecutor != null) { |
||||
|
dispatcherExecutor.shutdownNow(); |
||||
|
} |
||||
|
if (callbackExecutor != null) { |
||||
|
callbackExecutor.shutdownNow(); |
||||
|
} |
||||
|
if (timeoutExecutor != null) { |
||||
|
timeoutExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected abstract SettableFuture<V> create(); |
||||
|
|
||||
|
protected abstract F wrap(T task, SettableFuture<V> future); |
||||
|
|
||||
|
protected abstract ListenableFuture<V> execute(AsyncTaskContext<T, V> taskCtx); |
||||
|
|
||||
|
private void dispatch() { |
||||
|
log.info("Buffered rate executor thread started"); |
||||
|
while (!Thread.interrupted()) { |
||||
|
int curLvl = concurrencyLevel.get(); |
||||
|
AsyncTaskContext<T, V> taskCtx = null; |
||||
|
try { |
||||
|
if (curLvl <= concurrencyLimit) { |
||||
|
taskCtx = queue.take(); |
||||
|
final AsyncTaskContext<T, V> finalTaskCtx = taskCtx; |
||||
|
logTask("Processing", finalTaskCtx); |
||||
|
concurrencyLevel.incrementAndGet(); |
||||
|
long timeout = finalTaskCtx.getCreateTime() + maxWaitTime - System.currentTimeMillis(); |
||||
|
if (timeout > 0) { |
||||
|
totalLaunched.incrementAndGet(); |
||||
|
ListenableFuture<V> result = execute(finalTaskCtx); |
||||
|
result = Futures.withTimeout(result, timeout, TimeUnit.MILLISECONDS, timeoutExecutor); |
||||
|
Futures.addCallback(result, new FutureCallback<V>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable V result) { |
||||
|
logTask("Releasing", finalTaskCtx); |
||||
|
totalReleased.incrementAndGet(); |
||||
|
concurrencyLevel.decrementAndGet(); |
||||
|
finalTaskCtx.getFuture().set(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Throwable t) { |
||||
|
if (t instanceof TimeoutException) { |
||||
|
logTask("Expired During Execution", finalTaskCtx); |
||||
|
} else { |
||||
|
logTask("Failed", finalTaskCtx); |
||||
|
} |
||||
|
totalFailed.incrementAndGet(); |
||||
|
concurrencyLevel.decrementAndGet(); |
||||
|
finalTaskCtx.getFuture().setException(t); |
||||
|
log.debug("[{}] Failed to execute task: {}", finalTaskCtx.getId(), finalTaskCtx.getTask(), t); |
||||
|
} |
||||
|
}, callbackExecutor); |
||||
|
} else { |
||||
|
logTask("Expired Before Execution", finalTaskCtx); |
||||
|
totalExpired.incrementAndGet(); |
||||
|
concurrencyLevel.decrementAndGet(); |
||||
|
taskCtx.getFuture().setException(new TimeoutException()); |
||||
|
} |
||||
|
} else { |
||||
|
Thread.sleep(pollMs); |
||||
|
} |
||||
|
} catch (InterruptedException e) { |
||||
|
break; |
||||
|
} catch (Throwable e) { |
||||
|
if (taskCtx != null) { |
||||
|
log.debug("[{}] Failed to execute task: {}", taskCtx.getId(), taskCtx, e); |
||||
|
totalFailed.incrementAndGet(); |
||||
|
concurrencyLevel.decrementAndGet(); |
||||
|
} else { |
||||
|
log.debug("Failed to queue task:", e); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
log.info("Buffered rate executor thread stopped"); |
||||
|
} |
||||
|
|
||||
|
private void logTask(String action, AsyncTaskContext<T, V> taskCtx) { |
||||
|
if (log.isTraceEnabled()) { |
||||
|
log.trace("[{}] {} task: {}", taskCtx.getId(), action, taskCtx); |
||||
|
} else { |
||||
|
log.debug("[{}] {} task", taskCtx.getId(), action); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected int getQueueSize() { |
||||
|
return queue.size(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,22 @@ |
|||||
|
/** |
||||
|
* 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.util; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
public interface AsyncTask { |
||||
|
} |
||||
@ -0,0 +1,34 @@ |
|||||
|
/** |
||||
|
* 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.util; |
||||
|
|
||||
|
import com.google.common.util.concurrent.SettableFuture; |
||||
|
import lombok.Data; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
@Data |
||||
|
public class AsyncTaskContext<T extends AsyncTask, V> { |
||||
|
|
||||
|
private final UUID id; |
||||
|
private final T task; |
||||
|
private final SettableFuture<V> future; |
||||
|
private final long createTime; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,27 @@ |
|||||
|
/** |
||||
|
* 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.util; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 24.10.18. |
||||
|
*/ |
||||
|
public interface BufferedRateExecutor<T extends AsyncTask, F extends ListenableFuture> { |
||||
|
|
||||
|
F submit(T task); |
||||
|
|
||||
|
} |
||||
@ -1,182 +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.util; |
|
||||
|
|
||||
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.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.scheduling.annotation.Scheduled; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
import org.thingsboard.server.dao.exception.BufferLimitException; |
|
||||
|
|
||||
import java.util.concurrent.BlockingQueue; |
|
||||
import java.util.concurrent.CountDownLatch; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.LinkedBlockingQueue; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
|
|
||||
@Component |
|
||||
@Slf4j |
|
||||
@NoSqlAnyDao |
|
||||
public class BufferedRateLimiter implements AsyncRateLimiter { |
|
||||
|
|
||||
private final ListeningExecutorService pool = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10)); |
|
||||
|
|
||||
private final int permitsLimit; |
|
||||
private final int maxPermitWaitTime; |
|
||||
private final AtomicInteger permits; |
|
||||
private final BlockingQueue<LockedFuture> queue; |
|
||||
|
|
||||
private final AtomicInteger maxQueueSize = new AtomicInteger(); |
|
||||
private final AtomicInteger maxGrantedPermissions = new AtomicInteger(); |
|
||||
private final AtomicInteger totalGranted = new AtomicInteger(); |
|
||||
private final AtomicInteger totalReleased = new AtomicInteger(); |
|
||||
private final AtomicInteger totalRequested = new AtomicInteger(); |
|
||||
|
|
||||
public BufferedRateLimiter(@Value("${cassandra.query.buffer_size}") int queueLimit, |
|
||||
@Value("${cassandra.query.concurrent_limit}") int permitsLimit, |
|
||||
@Value("${cassandra.query.permit_max_wait_time}") int maxPermitWaitTime) { |
|
||||
this.permitsLimit = permitsLimit; |
|
||||
this.maxPermitWaitTime = maxPermitWaitTime; |
|
||||
this.permits = new AtomicInteger(); |
|
||||
this.queue = new LinkedBlockingQueue<>(queueLimit); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> acquireAsync() { |
|
||||
totalRequested.incrementAndGet(); |
|
||||
if (queue.isEmpty()) { |
|
||||
if (permits.incrementAndGet() <= permitsLimit) { |
|
||||
if (permits.get() > maxGrantedPermissions.get()) { |
|
||||
maxGrantedPermissions.set(permits.get()); |
|
||||
} |
|
||||
totalGranted.incrementAndGet(); |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
permits.decrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
return putInQueue(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void release() { |
|
||||
permits.decrementAndGet(); |
|
||||
totalReleased.incrementAndGet(); |
|
||||
reprocessQueue(); |
|
||||
} |
|
||||
|
|
||||
private void reprocessQueue() { |
|
||||
while (permits.get() < permitsLimit) { |
|
||||
if (permits.incrementAndGet() <= permitsLimit) { |
|
||||
if (permits.get() > maxGrantedPermissions.get()) { |
|
||||
maxGrantedPermissions.set(permits.get()); |
|
||||
} |
|
||||
LockedFuture lockedFuture = queue.poll(); |
|
||||
if (lockedFuture != null) { |
|
||||
totalGranted.incrementAndGet(); |
|
||||
lockedFuture.latch.countDown(); |
|
||||
} else { |
|
||||
permits.decrementAndGet(); |
|
||||
break; |
|
||||
} |
|
||||
} else { |
|
||||
permits.decrementAndGet(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private LockedFuture createLockedFuture() { |
|
||||
CountDownLatch latch = new CountDownLatch(1); |
|
||||
ListenableFuture<Void> future = pool.submit(() -> { |
|
||||
latch.await(); |
|
||||
return null; |
|
||||
}); |
|
||||
return new LockedFuture(latch, future, System.currentTimeMillis()); |
|
||||
} |
|
||||
|
|
||||
private ListenableFuture<Void> putInQueue() { |
|
||||
|
|
||||
int size = queue.size(); |
|
||||
if (size > maxQueueSize.get()) { |
|
||||
maxQueueSize.set(size); |
|
||||
} |
|
||||
|
|
||||
if (queue.remainingCapacity() > 0) { |
|
||||
try { |
|
||||
LockedFuture lockedFuture = createLockedFuture(); |
|
||||
if (!queue.offer(lockedFuture, 1, TimeUnit.SECONDS)) { |
|
||||
lockedFuture.cancelFuture(); |
|
||||
return Futures.immediateFailedFuture(new BufferLimitException()); |
|
||||
} |
|
||||
if(permits.get() < permitsLimit) { |
|
||||
reprocessQueue(); |
|
||||
} |
|
||||
if(permits.get() < permitsLimit) { |
|
||||
reprocessQueue(); |
|
||||
} |
|
||||
return lockedFuture.future; |
|
||||
} catch (InterruptedException e) { |
|
||||
return Futures.immediateFailedFuture(new BufferLimitException()); |
|
||||
} |
|
||||
} |
|
||||
return Futures.immediateFailedFuture(new BufferLimitException()); |
|
||||
} |
|
||||
|
|
||||
@Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") |
|
||||
public void printStats() { |
|
||||
int expiredCount = 0; |
|
||||
for (LockedFuture lockedFuture : queue) { |
|
||||
if (lockedFuture.isExpired()) { |
|
||||
lockedFuture.cancelFuture(); |
|
||||
expiredCount++; |
|
||||
} |
|
||||
} |
|
||||
log.info("Permits maxBuffer [{}] maxPermits [{}] expired [{}] currPermits [{}] currBuffer [{}] " + |
|
||||
"totalPermits [{}] totalRequests [{}] totalReleased [{}]", |
|
||||
maxQueueSize.getAndSet(0), maxGrantedPermissions.getAndSet(0), expiredCount, |
|
||||
permits.get(), queue.size(), |
|
||||
totalGranted.getAndSet(0), totalRequested.getAndSet(0), totalReleased.getAndSet(0)); |
|
||||
} |
|
||||
|
|
||||
private class LockedFuture { |
|
||||
final CountDownLatch latch; |
|
||||
final ListenableFuture<Void> future; |
|
||||
final long createTime; |
|
||||
|
|
||||
public LockedFuture(CountDownLatch latch, ListenableFuture<Void> future, long createTime) { |
|
||||
this.latch = latch; |
|
||||
this.future = future; |
|
||||
this.createTime = createTime; |
|
||||
} |
|
||||
|
|
||||
void cancelFuture() { |
|
||||
future.cancel(false); |
|
||||
latch.countDown(); |
|
||||
} |
|
||||
|
|
||||
boolean isExpired() { |
|
||||
return (System.currentTimeMillis() - createTime) > maxPermitWaitTime; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -1,142 +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.util; |
|
||||
|
|
||||
import com.google.common.util.concurrent.FutureCallback; |
|
||||
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 org.junit.Test; |
|
||||
import org.thingsboard.server.dao.exception.BufferLimitException; |
|
||||
|
|
||||
import javax.annotation.Nullable; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
|
|
||||
import static org.junit.Assert.assertEquals; |
|
||||
import static org.junit.Assert.assertFalse; |
|
||||
import static org.junit.Assert.assertTrue; |
|
||||
import static org.junit.Assert.fail; |
|
||||
|
|
||||
|
|
||||
public class BufferedRateLimiterTest { |
|
||||
|
|
||||
@Test |
|
||||
public void finishedFutureReturnedIfPermitsAreGranted() { |
|
||||
BufferedRateLimiter limiter = new BufferedRateLimiter(10, 10, 100); |
|
||||
ListenableFuture<Void> actual = limiter.acquireAsync(); |
|
||||
assertTrue(actual.isDone()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void notFinishedFutureReturnedIfPermitsAreNotGranted() { |
|
||||
BufferedRateLimiter limiter = new BufferedRateLimiter(10, 1, 100); |
|
||||
ListenableFuture<Void> actual1 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual2 = limiter.acquireAsync(); |
|
||||
assertTrue(actual1.isDone()); |
|
||||
assertFalse(actual2.isDone()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void failedFutureReturnedIfQueueIsfull() { |
|
||||
BufferedRateLimiter limiter = new BufferedRateLimiter(1, 1, 100); |
|
||||
ListenableFuture<Void> actual1 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual2 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual3 = limiter.acquireAsync(); |
|
||||
|
|
||||
assertTrue(actual1.isDone()); |
|
||||
assertFalse(actual2.isDone()); |
|
||||
assertTrue(actual3.isDone()); |
|
||||
try { |
|
||||
actual3.get(); |
|
||||
fail(); |
|
||||
} catch (Exception e) { |
|
||||
assertTrue(e instanceof ExecutionException); |
|
||||
Throwable actualCause = e.getCause(); |
|
||||
assertTrue(actualCause instanceof BufferLimitException); |
|
||||
assertEquals("Rate Limit Buffer is full", actualCause.getMessage()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void releasedPermitTriggerTasksFromQueue() throws InterruptedException { |
|
||||
BufferedRateLimiter limiter = new BufferedRateLimiter(10, 2, 100); |
|
||||
ListenableFuture<Void> actual1 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual2 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual3 = limiter.acquireAsync(); |
|
||||
ListenableFuture<Void> actual4 = limiter.acquireAsync(); |
|
||||
assertTrue(actual1.isDone()); |
|
||||
assertTrue(actual2.isDone()); |
|
||||
assertFalse(actual3.isDone()); |
|
||||
assertFalse(actual4.isDone()); |
|
||||
limiter.release(); |
|
||||
TimeUnit.MILLISECONDS.sleep(100L); |
|
||||
assertTrue(actual3.isDone()); |
|
||||
assertFalse(actual4.isDone()); |
|
||||
limiter.release(); |
|
||||
TimeUnit.MILLISECONDS.sleep(100L); |
|
||||
assertTrue(actual4.isDone()); |
|
||||
} |
|
||||
|
|
||||
@Test |
|
||||
public void permitsReleasedInConcurrentMode() throws InterruptedException { |
|
||||
BufferedRateLimiter limiter = new BufferedRateLimiter(10, 2, 100); |
|
||||
AtomicInteger actualReleased = new AtomicInteger(); |
|
||||
AtomicInteger actualRejected = new AtomicInteger(); |
|
||||
ListeningExecutorService pool = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(5)); |
|
||||
for (int i = 0; i < 100; i++) { |
|
||||
ListenableFuture<ListenableFuture<Void>> submit = pool.submit(limiter::acquireAsync); |
|
||||
Futures.addCallback(submit, new FutureCallback<ListenableFuture<Void>>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable ListenableFuture<Void> result) { |
|
||||
Futures.addCallback(result, new FutureCallback<Void>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable Void result) { |
|
||||
try { |
|
||||
TimeUnit.MILLISECONDS.sleep(100); |
|
||||
} catch (InterruptedException e) { |
|
||||
e.printStackTrace(); |
|
||||
} |
|
||||
limiter.release(); |
|
||||
actualReleased.incrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
actualRejected.incrementAndGet(); |
|
||||
} |
|
||||
}); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
} |
|
||||
}); |
|
||||
} |
|
||||
|
|
||||
TimeUnit.SECONDS.sleep(2); |
|
||||
assertTrue("Unexpected released count " + actualReleased.get(), |
|
||||
actualReleased.get() > 10 && actualReleased.get() < 20); |
|
||||
assertTrue("Unexpected rejected count " + actualRejected.get(), |
|
||||
actualRejected.get() > 80 && actualRejected.get() < 90); |
|
||||
|
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -1,13 +1,18 @@ |
|||||
# cassandra environment variables |
|
||||
CASSANDRA_DATA_DIR=/home/docker/cassandra_volume |
|
||||
|
|
||||
# postgres environment variables |
DOCKER_REPO=thingsboard |
||||
POSTGRES_DATA_DIR=/home/docker/postgres_volume |
|
||||
POSTGRES_DB=thingsboard |
|
||||
|
|
||||
# hsqldb environment variables |
JS_EXECUTOR_DOCKER_NAME=tb-js-executor |
||||
HSQLDB_DATA_DIR=/home/docker/hsqldb_volume |
TB_NODE_DOCKER_NAME=tb-node |
||||
|
WEB_UI_DOCKER_NAME=tb-web-ui |
||||
|
MQTT_TRANSPORT_DOCKER_NAME=tb-mqtt-transport |
||||
|
HTTP_TRANSPORT_DOCKER_NAME=tb-http-transport |
||||
|
COAP_TRANSPORT_DOCKER_NAME=tb-coap-transport |
||||
|
|
||||
# environment variables for schema init and insert system and demo data |
TB_VERSION=latest |
||||
ADD_SCHEMA_AND_SYSTEM_DATA=false |
|
||||
ADD_DEMO_DATA=false |
# Database used by ThingsBoard, can be either postgres (PostgreSQL) or cassandra (Cassandra). |
||||
|
# According to the database type corresponding docker service will be deployed (see docker-compose.postgres.yml, docker-compose.cassandra.yml for details). |
||||
|
|
||||
|
DATABASE=postgres |
||||
|
|
||||
|
KAFKA_TOPICS="js.eval.requests:100:1:delete --config=retention.ms=60000 --config=segment.bytes=26214400 --config=retention.bytes=104857600,tb.transport.api.requests:30:1:delete --config=retention.ms=60000 --config=segment.bytes=26214400 --config=retention.bytes=104857600,tb.rule-engine:30:1" |
||||
|
|||||
@ -0,0 +1,95 @@ |
|||||
|
# Docker configuration for ThingsBoard Microservices |
||||
|
|
||||
|
This folder containing scripts and Docker Compose configurations to run ThingsBoard in Microservices mode. |
||||
|
|
||||
|
## Prerequisites |
||||
|
|
||||
|
ThingsBoard Microservices are running in dockerized environment. |
||||
|
Before starting please make sure [Docker CE](https://docs.docker.com/install/) and [Docker Compose](https://docs.docker.com/compose/install/) are installed in your system. |
||||
|
|
||||
|
## Installation |
||||
|
|
||||
|
Before performing initial installation you can configure the type of database to be used with ThinsBoard. |
||||
|
In order to set database type change the value of `DATABASE` variable in `.env` file to one of the following: |
||||
|
|
||||
|
- `postgres` - use PostgreSQL database; |
||||
|
- `cassandra` - use Cassandra database; |
||||
|
|
||||
|
**NOTE**: According to the database type corresponding docker service will be deployed (see `docker-compose.postgres.yml`, `docker-compose.cassandra.yml` for details). |
||||
|
|
||||
|
Execute the following command to run installation: |
||||
|
|
||||
|
` |
||||
|
$ ./docker-install-tb.sh --loadDemo |
||||
|
` |
||||
|
|
||||
|
Where: |
||||
|
|
||||
|
- `--loadDemo` - optional argument. Whether to load additional demo data. |
||||
|
|
||||
|
## Running |
||||
|
|
||||
|
Execute the following command to start services: |
||||
|
|
||||
|
` |
||||
|
$ ./docker-start-services.sh |
||||
|
` |
||||
|
|
||||
|
After a while when all services will be successfully started you can open `http://{your-host-ip}` in you browser (for ex. `http://localhost`). |
||||
|
You should see ThingsBoard login page. |
||||
|
|
||||
|
Use the following default credentials: |
||||
|
|
||||
|
- **Systen Administrator**: sysadmin@thingsboard.org / sysadmin |
||||
|
|
||||
|
If you installed DataBase with demo data (using `--loadDemo` flag) you can also use the following credentials: |
||||
|
|
||||
|
- **Tenant Administrator**: tenant@thingsboard.org / tenant |
||||
|
- **Customer User**: customer@thingsboard.org / customer |
||||
|
|
||||
|
In case of any issues you can examine service logs for errors. |
||||
|
For example to see ThingsBoard node logs execute the following command: |
||||
|
|
||||
|
` |
||||
|
$ docker-compose logs -f tb1 |
||||
|
` |
||||
|
|
||||
|
Or use `docker-compose ps` to see the state of all the containers. |
||||
|
Use `docker-compose logs --f` to inspect the logs of all running services. |
||||
|
See [docker-compose logs](https://docs.docker.com/compose/reference/logs/) command reference for details. |
||||
|
|
||||
|
Execute the following command to stop services: |
||||
|
|
||||
|
` |
||||
|
$ ./docker-stop-services.sh |
||||
|
` |
||||
|
|
||||
|
Execute the following command to stop and completely remove deployed docker containers: |
||||
|
|
||||
|
` |
||||
|
$ ./docker-remove-services.sh |
||||
|
` |
||||
|
|
||||
|
Execute the following command to update particular or all services (pull newer docker image and rebuild container): |
||||
|
|
||||
|
` |
||||
|
$ ./docker-update-service.sh [SERVICE...] |
||||
|
` |
||||
|
|
||||
|
Where: |
||||
|
|
||||
|
- `[SERVICE...]` - list of services to update (defined in docker-compose configurations). If not specified all services will be updated. |
||||
|
|
||||
|
## Upgrading |
||||
|
|
||||
|
In case when database upgrade is needed, execute the following commands: |
||||
|
|
||||
|
``` |
||||
|
$ ./docker-stop-services.sh |
||||
|
$ ./docker-upgrade-tb.sh --fromVersion=[FROM_VERSION] |
||||
|
$ ./docker-start-services.sh |
||||
|
``` |
||||
|
|
||||
|
Where: |
||||
|
|
||||
|
- `FROM_VERSION` - from which version upgrade should be started. See [Upgrade Instructions](https://thingsboard.io/docs/user-guide/install/upgrade-instructions) for valid `fromVersion` values. |
||||
@ -1,24 +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. |
|
||||
# |
|
||||
|
|
||||
FROM openjdk:8-jre |
|
||||
|
|
||||
ADD install.sh /install.sh |
|
||||
ADD thingsboard.deb /thingsboard.deb |
|
||||
|
|
||||
RUN apt-get update \ |
|
||||
&& apt-get install -y nmap \ |
|
||||
&& chmod +x /install.sh |
|
||||
@ -1,12 +0,0 @@ |
|||||
VERSION=2.1.0 |
|
||||
PROJECT=thingsboard |
|
||||
APP=cassandra-setup |
|
||||
|
|
||||
build: |
|
||||
cp ../../application/target/thingsboard.deb . |
|
||||
docker build --pull -t ${PROJECT}/${APP}:${VERSION} -t ${PROJECT}/${APP}:latest . |
|
||||
rm thingsboard.deb |
|
||||
|
|
||||
push: build |
|
||||
docker push ${PROJECT}/${APP}:${VERSION} |
|
||||
docker push ${PROJECT}/${APP}:latest |
|
||||
@ -1,24 +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. |
|
||||
# |
|
||||
|
|
||||
FROM openjdk:8-jre |
|
||||
|
|
||||
ADD upgrade.sh /upgrade.sh |
|
||||
ADD thingsboard.deb /thingsboard.deb |
|
||||
|
|
||||
RUN apt-get update \ |
|
||||
&& apt-get install -y nmap \ |
|
||||
&& chmod +x /upgrade.sh |
|
||||
@ -1,12 +0,0 @@ |
|||||
VERSION=2.1.0 |
|
||||
PROJECT=thingsboard |
|
||||
APP=cassandra-upgrade |
|
||||
|
|
||||
build: |
|
||||
cp ../../application/target/thingsboard.deb . |
|
||||
docker build --pull -t ${PROJECT}/${APP}:${VERSION} -t ${PROJECT}/${APP}:latest . |
|
||||
rm thingsboard.deb |
|
||||
|
|
||||
push: build |
|
||||
docker push ${PROJECT}/${APP}:${VERSION} |
|
||||
docker push ${PROJECT}/${APP}:latest |
|
||||
@ -1,10 +0,0 @@ |
|||||
VERSION=2.1.0 |
|
||||
PROJECT=thingsboard |
|
||||
APP=cassandra |
|
||||
|
|
||||
build: |
|
||||
docker build --pull -t ${PROJECT}/${APP}:${VERSION} -t ${PROJECT}/${APP}:latest . |
|
||||
|
|
||||
push: build |
|
||||
docker push ${PROJECT}/${APP}:${VERSION} |
|
||||
docker push ${PROJECT}/${APP}:latest |
|
||||
@ -1,28 +0,0 @@ |
|||||
#!/usr/bin/env bash |
|
||||
# |
|
||||
# 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. |
|
||||
# |
|
||||
|
|
||||
if [[ $(nodetool status | grep $POD_IP) == *"UN"* ]]; then |
|
||||
if [[ $DEBUG ]]; then |
|
||||
echo "UN"; |
|
||||
fi |
|
||||
exit 0; |
|
||||
else |
|
||||
if [[ $DEBUG ]]; then |
|
||||
echo "Not Up"; |
|
||||
fi |
|
||||
exit 1; |
|
||||
fi |
|
||||
@ -0,0 +1,50 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
function additionalComposeArgs() { |
||||
|
source .env |
||||
|
ADDITIONAL_COMPOSE_ARGS="" |
||||
|
case $DATABASE in |
||||
|
postgres) |
||||
|
ADDITIONAL_COMPOSE_ARGS="-f docker-compose.postgres.yml" |
||||
|
;; |
||||
|
cassandra) |
||||
|
ADDITIONAL_COMPOSE_ARGS="-f docker-compose.cassandra.yml" |
||||
|
;; |
||||
|
*) |
||||
|
echo "Unknown DATABASE value specified: '${DATABASE}'. Should be either postgres or cassandra." >&2 |
||||
|
exit 1 |
||||
|
esac |
||||
|
echo $ADDITIONAL_COMPOSE_ARGS |
||||
|
} |
||||
|
|
||||
|
function additionalStartupServices() { |
||||
|
source .env |
||||
|
ADDITIONAL_STARTUP_SERVICES="" |
||||
|
case $DATABASE in |
||||
|
postgres) |
||||
|
ADDITIONAL_STARTUP_SERVICES=postgres |
||||
|
;; |
||||
|
cassandra) |
||||
|
ADDITIONAL_STARTUP_SERVICES=cassandra |
||||
|
;; |
||||
|
*) |
||||
|
echo "Unknown DATABASE value specified: '${DATABASE}'. Should be either postgres or cassandra." >&2 |
||||
|
exit 1 |
||||
|
esac |
||||
|
echo $ADDITIONAL_STARTUP_SERVICES |
||||
|
} |
||||
@ -1,27 +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. |
|
||||
# |
|
||||
|
|
||||
version: '3.3' |
|
||||
services: |
|
||||
redis: |
|
||||
image: redis:4.0 |
|
||||
networks: |
|
||||
- core |
|
||||
ports: |
|
||||
- "6379:6379" |
|
||||
|
|
||||
networks: |
|
||||
core: |
|
||||
@ -0,0 +1,24 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
set -e |
||||
|
|
||||
|
source compose-utils.sh |
||||
|
|
||||
|
ADDITIONAL_COMPOSE_ARGS=$(additionalComposeArgs) || exit $? |
||||
|
|
||||
|
docker-compose -f docker-compose.yml $ADDITIONAL_COMPOSE_ARGS down -v |
||||
@ -0,0 +1,26 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
./check-dirs.sh |
||||
|
|
||||
|
set -e |
||||
|
|
||||
|
source compose-utils.sh |
||||
|
|
||||
|
ADDITIONAL_COMPOSE_ARGS=$(additionalComposeArgs) || exit $? |
||||
|
|
||||
|
docker-compose -f docker-compose.yml $ADDITIONAL_COMPOSE_ARGS up -d |
||||
@ -0,0 +1,24 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
set -e |
||||
|
|
||||
|
source compose-utils.sh |
||||
|
|
||||
|
ADDITIONAL_COMPOSE_ARGS=$(additionalComposeArgs) || exit $? |
||||
|
|
||||
|
docker-compose -f docker-compose.yml $ADDITIONAL_COMPOSE_ARGS stop |
||||
@ -0,0 +1,25 @@ |
|||||
|
#!/bin/bash |
||||
|
# |
||||
|
# 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. |
||||
|
# |
||||
|
|
||||
|
set -e |
||||
|
|
||||
|
source compose-utils.sh |
||||
|
|
||||
|
ADDITIONAL_COMPOSE_ARGS=$(additionalComposeArgs) || exit $? |
||||
|
|
||||
|
docker-compose -f docker-compose.yml $ADDITIONAL_COMPOSE_ARGS pull $@ |
||||
|
docker-compose -f docker-compose.yml $ADDITIONAL_COMPOSE_ARGS up -d --no-deps --build $@ |
||||
@ -1,43 +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. |
|
||||
# |
|
||||
|
|
||||
apiVersion: v1 |
|
||||
kind: Pod |
|
||||
metadata: |
|
||||
name: cassandra-setup |
|
||||
spec: |
|
||||
containers: |
|
||||
- name: cassandra-setup |
|
||||
imagePullPolicy: Always |
|
||||
image: thingsboard/cassandra-setup:2.1.0 |
|
||||
env: |
|
||||
- name: ADD_DEMO_DATA |
|
||||
value: "true" |
|
||||
- name : CASSANDRA_HOST |
|
||||
value: "cassandra-headless" |
|
||||
- name : CASSANDRA_PORT |
|
||||
value: "9042" |
|
||||
- name : DATABASE_ENTITIES_TYPE |
|
||||
value: "cassandra" |
|
||||
- name : DATABASE_TS_TYPE |
|
||||
value: "cassandra" |
|
||||
- name : CASSANDRA_URL |
|
||||
value: "cassandra-headless:9042" |
|
||||
command: |
|
||||
- sh |
|
||||
- -c |
|
||||
- /install.sh |
|
||||
restartPolicy: Never |
|
||||
@ -1,45 +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. |
|
||||
# |
|
||||
|
|
||||
apiVersion: v1 |
|
||||
kind: Pod |
|
||||
metadata: |
|
||||
name: cassandra-upgrade |
|
||||
spec: |
|
||||
containers: |
|
||||
- name: cassandra-upgrade |
|
||||
imagePullPolicy: Always |
|
||||
image: thingsboard/cassandra-upgrade:2.1.0 |
|
||||
env: |
|
||||
- name: ADD_DEMO_DATA |
|
||||
value: "true" |
|
||||
- name : CASSANDRA_HOST |
|
||||
value: "cassandra-headless" |
|
||||
- name : CASSANDRA_PORT |
|
||||
value: "9042" |
|
||||
- name : DATABASE_ENTITIES_TYPE |
|
||||
value: "cassandra" |
|
||||
- name : DATABASE_TS_TYPE |
|
||||
value: "cassandra" |
|
||||
- name : CASSANDRA_URL |
|
||||
value: "cassandra-headless:9042" |
|
||||
- name : UPGRADE_FROM_VERSION |
|
||||
value: "1.4.0" |
|
||||
command: |
|
||||
- sh |
|
||||
- -c |
|
||||
- /upgrade.sh |
|
||||
restartPolicy: Never |
|
||||
@ -1,132 +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. |
|
||||
# |
|
||||
|
|
||||
apiVersion: v1 |
|
||||
kind: Service |
|
||||
metadata: |
|
||||
name: cassandra-headless |
|
||||
labels: |
|
||||
app: cassandra-headless |
|
||||
spec: |
|
||||
ports: |
|
||||
- port: 9042 |
|
||||
name: cql |
|
||||
clusterIP: None |
|
||||
selector: |
|
||||
app: cassandra |
|
||||
--- |
|
||||
apiVersion: "apps/v1beta1" |
|
||||
kind: StatefulSet |
|
||||
metadata: |
|
||||
name: cassandra |
|
||||
spec: |
|
||||
serviceName: cassandra-headless |
|
||||
replicas: 2 |
|
||||
template: |
|
||||
metadata: |
|
||||
labels: |
|
||||
app: cassandra |
|
||||
spec: |
|
||||
nodeSelector: |
|
||||
machinetype: other |
|
||||
affinity: |
|
||||
podAntiAffinity: |
|
||||
requiredDuringSchedulingIgnoredDuringExecution: |
|
||||
- labelSelector: |
|
||||
matchExpressions: |
|
||||
- key: "app" |
|
||||
operator: In |
|
||||
values: |
|
||||
- cassandra-headless |
|
||||
topologyKey: "kubernetes.io/hostname" |
|
||||
containers: |
|
||||
- name: cassandra |
|
||||
image: thingsboard/cassandra:2.1.0 |
|
||||
imagePullPolicy: Always |
|
||||
ports: |
|
||||
- containerPort: 7000 |
|
||||
name: intra-node |
|
||||
- containerPort: 7001 |
|
||||
name: tls-intra-node |
|
||||
- containerPort: 7199 |
|
||||
name: jmx |
|
||||
- containerPort: 9042 |
|
||||
name: cql |
|
||||
- containerPort: 9160 |
|
||||
name: thrift |
|
||||
securityContext: |
|
||||
capabilities: |
|
||||
add: |
|
||||
- IPC_LOCK |
|
||||
lifecycle: |
|
||||
preStop: |
|
||||
exec: |
|
||||
command: ["/bin/sh", "-c", "PID=$(pidof java) && kill $PID && while ps -p $PID > /dev/null; do sleep 1; done"] |
|
||||
env: |
|
||||
- name: MAX_HEAP_SIZE |
|
||||
value: 2048M |
|
||||
- name: HEAP_NEWSIZE |
|
||||
value: 100M |
|
||||
- name: CASSANDRA_SEEDS |
|
||||
value: "cassandra-0.cassandra-headless.default.svc.cluster.local" |
|
||||
- name: CASSANDRA_CLUSTER_NAME |
|
||||
value: "Thingsboard-Cluster" |
|
||||
- name: CASSANDRA_DC |
|
||||
value: "DC1-Thingsboard-Cluster" |
|
||||
- name: CASSANDRA_RACK |
|
||||
value: "Rack-Thingsboard-Cluster" |
|
||||
- name: CASSANDRA_AUTO_BOOTSTRAP |
|
||||
value: "false" |
|
||||
- name: POD_IP |
|
||||
valueFrom: |
|
||||
fieldRef: |
|
||||
fieldPath: status.podIP |
|
||||
- name: POD_NAMESPACE |
|
||||
valueFrom: |
|
||||
fieldRef: |
|
||||
fieldPath: metadata.namespace |
|
||||
readinessProbe: |
|
||||
exec: |
|
||||
command: |
|
||||
- /bin/bash |
|
||||
- -c |
|
||||
- /ready-probe.sh |
|
||||
initialDelaySeconds: 15 |
|
||||
timeoutSeconds: 5 |
|
||||
volumeMounts: |
|
||||
- name: cassandra-data |
|
||||
mountPath: /var/lib/cassandra/data |
|
||||
- name: cassandra-commitlog |
|
||||
mountPath: /var/lib/cassandra/commitlog |
|
||||
volumeClaimTemplates: |
|
||||
- metadata: |
|
||||
name: cassandra-data |
|
||||
annotations: |
|
||||
volume.beta.kubernetes.io/storage-class: fast |
|
||||
spec: |
|
||||
accessModes: [ "ReadWriteOnce" ] |
|
||||
resources: |
|
||||
requests: |
|
||||
storage: 3Gi |
|
||||
- metadata: |
|
||||
name: cassandra-commitlog |
|
||||
annotations: |
|
||||
volume.beta.kubernetes.io/storage-class: fast |
|
||||
spec: |
|
||||
accessModes: [ "ReadWriteOnce" ] |
|
||||
resources: |
|
||||
requests: |
|
||||
storage: 2Gi |
|
||||
@ -1,33 +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. |
|
||||
# |
|
||||
|
|
||||
--- |
|
||||
apiVersion: storage.k8s.io/v1beta1 |
|
||||
kind: StorageClass |
|
||||
metadata: |
|
||||
name: slow |
|
||||
provisioner: kubernetes.io/gce-pd |
|
||||
parameters: |
|
||||
type: pd-standard |
|
||||
--- |
|
||||
apiVersion: storage.k8s.io/v1beta1 |
|
||||
kind: StorageClass |
|
||||
metadata: |
|
||||
name: fast |
|
||||
provisioner: kubernetes.io/gce-pd |
|
||||
parameters: |
|
||||
type: pd-ssd |
|
||||
--- |
|
||||
@ -1,93 +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. |
|
||||
# |
|
||||
|
|
||||
--- |
|
||||
apiVersion: v1 |
|
||||
kind: Service |
|
||||
metadata: |
|
||||
labels: |
|
||||
name: redis-service |
|
||||
name: redis-service |
|
||||
spec: |
|
||||
ports: |
|
||||
- name: redis-service |
|
||||
protocol: TCP |
|
||||
port: 6379 |
|
||||
targetPort: 6379 |
|
||||
selector: |
|
||||
app: redis |
|
||||
--- |
|
||||
apiVersion: v1 |
|
||||
kind: ConfigMap |
|
||||
metadata: |
|
||||
name: redis-conf |
|
||||
data: |
|
||||
redis.conf: | |
|
||||
appendonly yes |
|
||||
protected-mode no |
|
||||
bind 0.0.0.0 |
|
||||
port 6379 |
|
||||
dir /var/lib/redis |
|
||||
--- |
|
||||
apiVersion: apps/v1beta1 |
|
||||
kind: StatefulSet |
|
||||
metadata: |
|
||||
name: redis |
|
||||
spec: |
|
||||
serviceName: redis-service |
|
||||
replicas: 1 |
|
||||
template: |
|
||||
metadata: |
|
||||
labels: |
|
||||
app: redis |
|
||||
spec: |
|
||||
terminationGracePeriodSeconds: 10 |
|
||||
containers: |
|
||||
- name: redis |
|
||||
image: redis:4.0.9 |
|
||||
command: |
|
||||
- redis-server |
|
||||
args: |
|
||||
- /etc/redis/redis.conf |
|
||||
resources: |
|
||||
requests: |
|
||||
cpu: 100m |
|
||||
memory: 100Mi |
|
||||
ports: |
|
||||
- containerPort: 6379 |
|
||||
name: redis |
|
||||
volumeMounts: |
|
||||
- name: redis-data |
|
||||
mountPath: /var/lib/redis |
|
||||
- name: redis-conf |
|
||||
mountPath: /etc/redis |
|
||||
volumes: |
|
||||
- name: redis-conf |
|
||||
configMap: |
|
||||
name: redis-conf |
|
||||
items: |
|
||||
- key: redis.conf |
|
||||
path: redis.conf |
|
||||
volumeClaimTemplates: |
|
||||
- metadata: |
|
||||
name: redis-data |
|
||||
annotations: |
|
||||
volume.beta.kubernetes.io/storage-class: fast |
|
||||
spec: |
|
||||
accessModes: [ "ReadWriteOnce" ] |
|
||||
resources: |
|
||||
requests: |
|
||||
storage: 1Gi |
|
||||
@ -1,156 +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. |
|
||||
# |
|
||||
|
|
||||
--- |
|
||||
apiVersion: v1 |
|
||||
kind: Service |
|
||||
metadata: |
|
||||
name: tb-service |
|
||||
labels: |
|
||||
app: tb-service |
|
||||
spec: |
|
||||
ports: |
|
||||
- port: 8080 |
|
||||
name: ui |
|
||||
- port: 1883 |
|
||||
name: mqtt |
|
||||
- port: 5683 |
|
||||
name: coap |
|
||||
selector: |
|
||||
app: tb |
|
||||
type: LoadBalancer |
|
||||
--- |
|
||||
apiVersion: policy/v1beta1 |
|
||||
kind: PodDisruptionBudget |
|
||||
metadata: |
|
||||
name: tb-budget |
|
||||
spec: |
|
||||
selector: |
|
||||
matchLabels: |
|
||||
app: tb |
|
||||
minAvailable: 3 |
|
||||
--- |
|
||||
apiVersion: v1 |
|
||||
kind: ConfigMap |
|
||||
metadata: |
|
||||
name: tb-config |
|
||||
data: |
|
||||
zookeeper.enabled: "true" |
|
||||
zookeeper.url: "zk-headless" |
|
||||
cassandra.url: "cassandra-headless:9042" |
|
||||
cassandra.host: "cassandra-headless" |
|
||||
cassandra.port: "9042" |
|
||||
database.type: "cassandra" |
|
||||
cache.type: "redis" |
|
||||
redis.host: "redis-service" |
|
||||
--- |
|
||||
apiVersion: apps/v1beta1 |
|
||||
kind: StatefulSet |
|
||||
metadata: |
|
||||
name: tb |
|
||||
spec: |
|
||||
serviceName: "tb-service" |
|
||||
replicas: 3 |
|
||||
template: |
|
||||
metadata: |
|
||||
labels: |
|
||||
app: tb |
|
||||
spec: |
|
||||
nodeSelector: |
|
||||
machinetype: tb |
|
||||
affinity: |
|
||||
podAntiAffinity: |
|
||||
requiredDuringSchedulingIgnoredDuringExecution: |
|
||||
- labelSelector: |
|
||||
matchExpressions: |
|
||||
- key: "app" |
|
||||
operator: In |
|
||||
values: |
|
||||
- tb-service |
|
||||
topologyKey: "kubernetes.io/hostname" |
|
||||
containers: |
|
||||
- name: tb |
|
||||
imagePullPolicy: Always |
|
||||
image: thingsboard/application:2.1.0 |
|
||||
ports: |
|
||||
- containerPort: 8080 |
|
||||
name: ui |
|
||||
- containerPort: 1883 |
|
||||
name: mqtt |
|
||||
- containerPort: 5683 |
|
||||
name: coap |
|
||||
- containerPort: 9001 |
|
||||
name: rpc |
|
||||
env: |
|
||||
- name: ZOOKEEPER_ENABLED |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: zookeeper.enabled |
|
||||
- name: ZOOKEEPER_URL |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: zookeeper.url |
|
||||
- name : CASSANDRA_HOST |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: cassandra.host |
|
||||
- name : CASSANDRA_PORT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: cassandra.port |
|
||||
- name : CASSANDRA_URL |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: cassandra.url |
|
||||
- name: DATABASE_ENTITIES_TYPE |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: database.type |
|
||||
- name: DATABASE_TS_TYPE |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: database.type |
|
||||
- name : RPC_HOST |
|
||||
valueFrom: |
|
||||
fieldRef: |
|
||||
fieldPath: status.podIP |
|
||||
- name: CACHE_TYPE |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: cache.type |
|
||||
- name: REDIS_HOST |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: tb-config |
|
||||
key: redis.host |
|
||||
command: |
|
||||
- sh |
|
||||
- -c |
|
||||
- /run-application.sh |
|
||||
livenessProbe: |
|
||||
httpGet: |
|
||||
path: /login |
|
||||
port: ui-port |
|
||||
initialDelaySeconds: 120 |
|
||||
timeoutSeconds: 10 |
|
||||
@ -1,190 +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. |
|
||||
# |
|
||||
|
|
||||
apiVersion: v1 |
|
||||
kind: Service |
|
||||
metadata: |
|
||||
name: zk-headless |
|
||||
labels: |
|
||||
app: zk-headless |
|
||||
spec: |
|
||||
ports: |
|
||||
- port: 2888 |
|
||||
name: server |
|
||||
- port: 3888 |
|
||||
name: leader-election |
|
||||
clusterIP: None |
|
||||
selector: |
|
||||
app: zk |
|
||||
--- |
|
||||
apiVersion: v1 |
|
||||
kind: ConfigMap |
|
||||
metadata: |
|
||||
name: zk-config |
|
||||
data: |
|
||||
ensemble: "zk-0;zk-1;zk-2" |
|
||||
replicas: "3" |
|
||||
jvm.heap: "500m" |
|
||||
tick: "2000" |
|
||||
init: "10" |
|
||||
sync: "5" |
|
||||
client.cnxns: "60" |
|
||||
snap.retain: "3" |
|
||||
purge.interval: "1" |
|
||||
client.port: "2181" |
|
||||
server.port: "2888" |
|
||||
election.port: "3888" |
|
||||
--- |
|
||||
apiVersion: policy/v1beta1 |
|
||||
kind: PodDisruptionBudget |
|
||||
metadata: |
|
||||
name: zk-budget |
|
||||
spec: |
|
||||
selector: |
|
||||
matchLabels: |
|
||||
app: zk |
|
||||
minAvailable: 3 |
|
||||
--- |
|
||||
apiVersion: apps/v1beta1 |
|
||||
kind: StatefulSet |
|
||||
metadata: |
|
||||
name: zk |
|
||||
spec: |
|
||||
serviceName: zk-headless |
|
||||
replicas: 3 |
|
||||
template: |
|
||||
metadata: |
|
||||
labels: |
|
||||
app: zk |
|
||||
annotations: |
|
||||
pod.alpha.kubernetes.io/initialized: "true" |
|
||||
spec: |
|
||||
nodeSelector: |
|
||||
machinetype: other |
|
||||
affinity: |
|
||||
podAntiAffinity: |
|
||||
requiredDuringSchedulingIgnoredDuringExecution: |
|
||||
- labelSelector: |
|
||||
matchExpressions: |
|
||||
- key: "app" |
|
||||
operator: In |
|
||||
values: |
|
||||
- zk-headless |
|
||||
topologyKey: "kubernetes.io/hostname" |
|
||||
containers: |
|
||||
- name: zk |
|
||||
imagePullPolicy: Always |
|
||||
image: thingsboard/zk:2.1.0 |
|
||||
ports: |
|
||||
- containerPort: 2181 |
|
||||
name: client |
|
||||
- containerPort: 2888 |
|
||||
name: server |
|
||||
- containerPort: 3888 |
|
||||
name: leader-election |
|
||||
env: |
|
||||
- name : ZK_ENSEMBLE |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: ensemble |
|
||||
- name : ZK_REPLICAS |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: replicas |
|
||||
- name : ZK_HEAP_SIZE |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: jvm.heap |
|
||||
- name : ZK_TICK_TIME |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: tick |
|
||||
- name : ZK_INIT_LIMIT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: init |
|
||||
- name : ZK_SYNC_LIMIT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: tick |
|
||||
- name : ZK_MAX_CLIENT_CNXNS |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: client.cnxns |
|
||||
- name: ZK_SNAP_RETAIN_COUNT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: snap.retain |
|
||||
- name: ZK_PURGE_INTERVAL |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: purge.interval |
|
||||
- name: ZK_CLIENT_PORT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: client.port |
|
||||
- name: ZK_SERVER_PORT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: server.port |
|
||||
- name: ZK_ELECTION_PORT |
|
||||
valueFrom: |
|
||||
configMapKeyRef: |
|
||||
name: zk-config |
|
||||
key: election.port |
|
||||
command: |
|
||||
- sh |
|
||||
- -c |
|
||||
- zk-gen-config.sh && zkServer.sh start-foreground |
|
||||
readinessProbe: |
|
||||
exec: |
|
||||
command: |
|
||||
- "zk-ok.sh" |
|
||||
initialDelaySeconds: 15 |
|
||||
timeoutSeconds: 5 |
|
||||
livenessProbe: |
|
||||
exec: |
|
||||
command: |
|
||||
- "zk-ok.sh" |
|
||||
initialDelaySeconds: 15 |
|
||||
timeoutSeconds: 5 |
|
||||
volumeMounts: |
|
||||
- name: zkdatadir |
|
||||
mountPath: /var/lib/zookeeper |
|
||||
securityContext: |
|
||||
runAsUser: 1000 |
|
||||
fsGroup: 1000 |
|
||||
volumeClaimTemplates: |
|
||||
- metadata: |
|
||||
name: zkdatadir |
|
||||
annotations: |
|
||||
volume.beta.kubernetes.io/storage-class: slow |
|
||||
spec: |
|
||||
accessModes: [ "ReadWriteOnce" ] |
|
||||
resources: |
|
||||
requests: |
|
||||
storage: 1Gi |
|
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue