Browse Source

Few performance tuning parameters

pull/782/head
Andrew Shvayka 9 years ago
parent
commit
b835d665c3
  1. 1
      application/src/main/resources/thingsboard.yml
  2. 8
      dao/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java
  3. 8
      transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java
  4. 5
      transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

1
application/src/main/resources/thingsboard.yml

@ -87,6 +87,7 @@ mqtt:
leak_detector_level: "${NETTY_LEASK_DETECTOR_LVL:DISABLED}"
boss_group_thread_count: "${NETTY_BOSS_GROUP_THREADS:1}"
worker_group_thread_count: "${NETTY_WORKER_GROUP_THREADS:12}"
max_payload_size: "${NETTY_MAX_PAYLOAD_SIZE:65536}"
# MQTT SSL configuration
ssl:
# Enable/disable SSL support

8
dao/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java

@ -60,6 +60,10 @@ public abstract class AbstractCassandraCluster {
private long initTimeout;
@Value("${cassandra.init_retry_interval_ms}")
private long initRetryInterval;
@Value("${cassandra.max_requests_per_connection_local:128}")
private int max_requests_local;
@Value("${cassandra.max_requests_per_connection_remote:128}")
private int max_requests_remote;
@Autowired
private CassandraSocketOptions socketOpts;
@ -90,8 +94,8 @@ public abstract class AbstractCassandraCluster {
.withClusterName(clusterName)
.withSocketOptions(socketOpts.getOpts())
.withPoolingOptions(new PoolingOptions()
.setMaxRequestsPerConnection(HostDistance.LOCAL, 32768)
.setMaxRequestsPerConnection(HostDistance.REMOTE, 32768));
.setMaxRequestsPerConnection(HostDistance.LOCAL, max_requests_local)
.setMaxRequestsPerConnection(HostDistance.REMOTE, max_requests_remote));
this.clusterBuilder.withQueryOptions(queryOpts.getOpts());
this.clusterBuilder.withCompression(StringUtils.isEmpty(compression) ? Compression.NONE : Compression.valueOf(compression.toUpperCase()));
if (ssl) {

8
transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java

@ -33,8 +33,6 @@ import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor;
*/
public class MqttTransportServerInitializer extends ChannelInitializer<SocketChannel> {
private static final int MAX_PAYLOAD_SIZE = 64 * 1024 * 1024;
private final SessionMsgProcessor processor;
private final DeviceService deviceService;
private final DeviceAuthService authService;
@ -42,10 +40,11 @@ public class MqttTransportServerInitializer extends ChannelInitializer<SocketCha
private final MqttTransportAdaptor adaptor;
private final MqttSslHandlerProvider sslHandlerProvider;
private final QuotaService quotaService;
private final int maxPayloadSize;
public MqttTransportServerInitializer(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, RelationService relationService,
MqttTransportAdaptor adaptor, MqttSslHandlerProvider sslHandlerProvider,
QuotaService quotaService) {
QuotaService quotaService, int maxPayloadSize) {
this.processor = processor;
this.deviceService = deviceService;
this.authService = authService;
@ -53,6 +52,7 @@ public class MqttTransportServerInitializer extends ChannelInitializer<SocketCha
this.adaptor = adaptor;
this.sslHandlerProvider = sslHandlerProvider;
this.quotaService = quotaService;
this.maxPayloadSize = maxPayloadSize;
}
@Override
@ -63,7 +63,7 @@ public class MqttTransportServerInitializer extends ChannelInitializer<SocketCha
sslHandler = sslHandlerProvider.getSslHandler();
pipeline.addLast(sslHandler);
}
pipeline.addLast("decoder", new MqttDecoder(MAX_PAYLOAD_SIZE));
pipeline.addLast("decoder", new MqttDecoder(maxPayloadSize));
pipeline.addLast("encoder", MqttEncoder.INSTANCE);
MqttTransportHandler handler = new MqttTransportHandler(processor, deviceService, authService, relationService,

5
transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

@ -82,7 +82,8 @@ public class MqttTransportService {
private Integer bossGroupThreadCount;
@Value("${mqtt.netty.worker_group_thread_count}")
private Integer workerGroupThreadCount;
@Value("${mqtt.netty.max_payload_size}")
private Integer maxPayloadSize;
private MqttTransportAdaptor adaptor;
@ -106,7 +107,7 @@ public class MqttTransportService {
b.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new MqttTransportServerInitializer(processor, deviceService, authService, relationService,
adaptor, sslHandlerProvider, quotaService));
adaptor, sslHandlerProvider, quotaService, maxPayloadSize));
serverChannel = b.bind(host, port).sync().channel();
log.info("Mqtt transport started!");

Loading…
Cancel
Save