diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index c370616be6..8f05422ea9 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -121,7 +121,7 @@ public class AppActor extends ContextAwareActor { private void broadcast(Object msg) { pluginManager.broadcast(msg); - tenantActors.values().stream().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); + tenantActors.values().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); } private void onToRuleMsg(ToRuleActorMsg msg) { diff --git a/application/src/main/java/org/thingsboard/server/actors/session/ASyncMsgProcessor.java b/application/src/main/java/org/thingsboard/server/actors/session/ASyncMsgProcessor.java index eb812dffdd..90b0cb316a 100644 --- a/application/src/main/java/org/thingsboard/server/actors/session/ASyncMsgProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/session/ASyncMsgProcessor.java @@ -111,7 +111,7 @@ class ASyncMsgProcessor extends AbstractSessionActorMsgProcessor { Optional newTargetServer = systemContext.getRoutingService().resolve(getDeviceId()); if (!newTargetServer.equals(currentTargetServer)) { currentTargetServer = newTargetServer; - pendingMap.values().stream().forEach(v -> { + pendingMap.values().forEach(v -> { forwardToAppActor(context, v, currentTargetServer); if (currentTargetServer.isPresent()) { logger.debug("[{}] Forwarded msg to new server: {}", sessionId, currentTargetServer.get()); diff --git a/application/src/main/java/org/thingsboard/server/actors/session/SessionManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/session/SessionManagerActor.java index 44eff1680e..c69946f1f6 100644 --- a/application/src/main/java/org/thingsboard/server/actors/session/SessionManagerActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/session/SessionManagerActor.java @@ -66,7 +66,7 @@ public class SessionManagerActor extends ContextAwareActor { } private void broadcast(Object msg) { - sessionActors.values().stream().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); + sessionActors.values().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); } private void onSessionTimeout(SessionTimeoutMsg msg) { diff --git a/application/src/main/java/org/thingsboard/server/actors/shared/plugin/PluginManager.java b/application/src/main/java/org/thingsboard/server/actors/shared/plugin/PluginManager.java index c581c41123..c0380d509f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/shared/plugin/PluginManager.java +++ b/application/src/main/java/org/thingsboard/server/actors/shared/plugin/PluginManager.java @@ -74,7 +74,7 @@ public abstract class PluginManager { } public void broadcast(Object msg) { - pluginActors.values().stream().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); + pluginActors.values().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); } public void remove(PluginId id) { diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index 965c652676..c8d5243f67 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -100,7 +100,7 @@ public class TenantActor extends ContextAwareActor { private void broadcast(Object msg) { pluginManager.broadcast(msg); - deviceActors.values().stream().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); + deviceActors.values().forEach(actorRef -> actorRef.tell(msg, ActorRef.noSender())); } private void onToDeviceActorMsg(ToDeviceActorMsg msg) { diff --git a/application/src/main/java/org/thingsboard/server/service/cluster/discovery/ZkDiscoveryService.java b/application/src/main/java/org/thingsboard/server/service/cluster/discovery/ZkDiscoveryService.java index 29e9b3cce4..dd784f0e34 100644 --- a/application/src/main/java/org/thingsboard/server/service/cluster/discovery/ZkDiscoveryService.java +++ b/application/src/main/java/org/thingsboard/server/service/cluster/discovery/ZkDiscoveryService.java @@ -166,7 +166,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi @Override public void onApplicationEvent(ApplicationReadyEvent applicationReadyEvent) { publishCurrentServer(); - getOtherServers().stream().forEach( + getOtherServers().forEach( server -> log.info("Found active server: [{}:{}]", server.getHost(), server.getPort()) ); } @@ -194,13 +194,13 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi log.info("Processing [{}] event for [{}:{}]", pathChildrenCacheEvent.getType(), instance.getHost(), instance.getPort()); switch (pathChildrenCacheEvent.getType()) { case CHILD_ADDED: - listeners.stream().forEach(listener -> listener.onServerAdded(instance)); + listeners.forEach(listener -> listener.onServerAdded(instance)); break; case CHILD_UPDATED: - listeners.stream().forEach(listener -> listener.onServerUpdated(instance)); + listeners.forEach(listener -> listener.onServerUpdated(instance)); break; case CHILD_REMOVED: - listeners.stream().forEach(listener -> listener.onServerRemoved(instance)); + listeners.forEach(listener -> listener.onServerRemoved(instance)); break; } } diff --git a/application/src/main/java/org/thingsboard/server/service/cluster/routing/ConsistentClusterRoutingService.java b/application/src/main/java/org/thingsboard/server/service/cluster/routing/ConsistentClusterRoutingService.java index 3c9ecf8a9b..7a3c7ac800 100644 --- a/application/src/main/java/org/thingsboard/server/service/cluster/routing/ConsistentClusterRoutingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cluster/routing/ConsistentClusterRoutingService.java @@ -135,7 +135,7 @@ public class ConsistentClusterRoutingService implements ClusterRoutingService, D private void logCircle() { log.trace("Consistent Hash Circle Start"); - circle.entrySet().stream().forEach((e) -> log.debug("{} -> {}", e.getKey(), e.getValue().getServerAddress())); + circle.entrySet().forEach((e) -> log.debug("{} -> {}", e.getKey(), e.getValue().getServerAddress())); log.trace("Consistent Hash Circle End"); } diff --git a/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java b/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java index a51464c4f4..975b52a5e7 100644 --- a/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java +++ b/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java @@ -31,7 +31,6 @@ import org.thingsboard.server.dao.component.ComponentDescriptorService; import org.thingsboard.server.extensions.api.component.*; import javax.annotation.PostConstruct; -import java.io.IOException; import java.lang.annotation.Annotation; import java.util.*; import java.util.stream.Collectors; @@ -72,7 +71,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe } private void registerComponents(Collection comps) { - comps.stream().forEach(c -> components.put(c.getClazz(), c)); + comps.forEach(c -> components.put(c.getClazz(), c)); } private List persist(Set filterDefs, ComponentType type) { @@ -119,7 +118,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe throw new RuntimeException("Plugin " + def.getBeanClassName() + "action " + actionClazz.getName() + " has wrong component type!"); } } - scannedComponent.setActions(Arrays.asList(pluginAnnotation.actions()).stream().map(action -> action.getName()).collect(Collectors.joining(","))); + scannedComponent.setActions(Arrays.stream(pluginAnnotation.actions()).map(action -> action.getName()).collect(Collectors.joining(","))); break; default: throw new RuntimeException(type + " is not supported yet!"); diff --git a/application/src/main/java/org/thingsboard/server/service/security/model/SecurityUser.java b/application/src/main/java/org/thingsboard/server/service/security/model/SecurityUser.java index 83b87abc7b..64569682f2 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/model/SecurityUser.java +++ b/application/src/main/java/org/thingsboard/server/service/security/model/SecurityUser.java @@ -20,9 +20,9 @@ import org.springframework.security.core.authority.SimpleGrantedAuthority; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.UserId; -import java.util.Arrays; import java.util.Collection; import java.util.stream.Collectors; +import java.util.stream.Stream; public class SecurityUser extends User { @@ -46,7 +46,7 @@ public class SecurityUser extends User { public Collection getAuthorities() { if (authorities == null) { - authorities = Arrays.asList(SecurityUser.this.getAuthority()).stream() + authorities = Stream.of(SecurityUser.this.getAuthority()) .map(authority -> new SimpleGrantedAuthority(authority.name())) .collect(Collectors.toList()); } diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java index 538e8f9ea9..20ceb77e4a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java @@ -129,8 +129,10 @@ public abstract class AbstractControllerTest { @Autowired void setConverters(HttpMessageConverter[] converters) { - this.mappingJackson2HttpMessageConverter = Arrays.asList(converters).stream().filter( - hmc -> hmc instanceof MappingJackson2HttpMessageConverter).findAny().get(); + this.mappingJackson2HttpMessageConverter = Arrays.stream(converters) + .filter(hmc -> hmc instanceof MappingJackson2HttpMessageConverter) + .findAny() + .get(); Assert.assertNotNull("the JSON message converter must not be null", this.mappingJackson2HttpMessageConverter); diff --git a/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java b/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java index 8d22343e07..db90b89a95 100644 --- a/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java @@ -61,8 +61,10 @@ public class AbstractFeatureIntegrationTest { @Autowired void setConverters(HttpMessageConverter[] converters) { - this.mappingJackson2HttpMessageConverter = Arrays.asList(converters).stream().filter( - hmc -> hmc instanceof MappingJackson2HttpMessageConverter).findAny().get(); + this.mappingJackson2HttpMessageConverter = Arrays.stream(converters) + .filter(hmc -> hmc instanceof MappingJackson2HttpMessageConverter) + .findAny() + .get(); assertNotNull("the JSON message converter must not be null", this.mappingJackson2HttpMessageConverter); diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java index 4c542e3bf2..ce4dd817bd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java @@ -140,7 +140,7 @@ public class BaseAttributesDao extends AbstractDao implements AttributesDao { List rows = resultSet.all(); List entries = new ArrayList<>(rows.size()); if (!rows.isEmpty()) { - rows.stream().forEach(row -> { + rows.forEach(row -> { String key = row.getString(ModelConstants.ATTRIBUTE_KEY_COLUMN); AttributeKvEntry kvEntry = convertResultToAttributesKvEntry(key, row); if (kvEntry != null) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java index 79134f7b0b..851c770439 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java @@ -143,7 +143,7 @@ public class BaseTimeseriesDao extends AbstractDao implements TimeseriesDao { public List convertResultToTsKvEntryList(List rows) { List entries = new ArrayList<>(rows.size()); if (!rows.isEmpty()) { - rows.stream().forEach(row -> { + rows.forEach(row -> { TsKvEntry kvEntry = convertResultToTsKvEntry(row); if (kvEntry != null) { entries.add(kvEntry); diff --git a/extensions-api/src/main/java/org/thingsboard/server/extensions/api/plugins/handlers/DefaultWebsocketMsgHandler.java b/extensions-api/src/main/java/org/thingsboard/server/extensions/api/plugins/handlers/DefaultWebsocketMsgHandler.java index fab11bbc07..1633c7f133 100644 --- a/extensions-api/src/main/java/org/thingsboard/server/extensions/api/plugins/handlers/DefaultWebsocketMsgHandler.java +++ b/extensions-api/src/main/java/org/thingsboard/server/extensions/api/plugins/handlers/DefaultWebsocketMsgHandler.java @@ -91,7 +91,7 @@ public class DefaultWebsocketMsgHandler implements WebsocketMsgHandler { } public void clear(PluginContext ctx) { - wsSessionsMap.values().stream().forEach(v -> { + wsSessionsMap.values().forEach(v -> { try { ctx.close(v.getSessionRef()); } catch (IOException e) { diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MethodNameFilter.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MethodNameFilter.java index 21180d791b..7123f3eb13 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MethodNameFilter.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MethodNameFilter.java @@ -40,7 +40,9 @@ public class MethodNameFilter extends SimpleRuleLifecycleComponent implements Ru @Override public void init(MethodNameFilterConfiguration configuration) { - methods = Arrays.asList(configuration.getMethodNames()).stream().map(m -> m.getName()).collect(Collectors.toSet()); + methods = Arrays.stream(configuration.getMethodNames()) + .map(m -> m.getName()) + .collect(Collectors.toSet()); } @Override diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MsgTypeFilter.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MsgTypeFilter.java index 84deea59db..737bee6c52 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MsgTypeFilter.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/filter/MsgTypeFilter.java @@ -39,7 +39,7 @@ public class MsgTypeFilter extends SimpleRuleLifecycleComponent implements RuleF @Override public void init(MsgTypeFilterConfiguration configuration) { - msgTypes = Arrays.asList(configuration.getMessageTypes()).stream().map(type -> { + msgTypes = Arrays.stream(configuration.getMessageTypes()).map(type -> { switch (type) { case "GET_ATTRIBUTES": return MsgType.GET_ATTRIBUTES_REQUEST; diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/mail/MailPlugin.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/mail/MailPlugin.java index 52fd2e9e8a..c8a7ad8361 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/mail/MailPlugin.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/mail/MailPlugin.java @@ -75,7 +75,7 @@ public class MailPlugin extends AbstractPlugin implemen if (configuration.getOtherProperties() != null) { Properties mailProperties = new Properties(); configuration.getOtherProperties() - .stream().forEach(p -> mailProperties.put(p.getKey(), p.getValue())); + .forEach(p -> mailProperties.put(p.getKey(), p.getValue())); mail.setJavaMailProperties(mailProperties); } mailSender = mail; diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/SubscriptionManager.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/SubscriptionManager.java index 637500e66f..190d9ffa0c 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/SubscriptionManager.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/SubscriptionManager.java @@ -68,7 +68,7 @@ public class SubscriptionManager { registerSubscription(sessionId, deviceId, subscription); List missedUpdates = new ArrayList<>(); if (subscription.getType() == SubscriptionType.ATTRIBUTES) { - subscription.getKeyStates().entrySet().stream().forEach(e -> { + subscription.getKeyStates().entrySet().forEach(e -> { Optional latestOpt = ctx.loadAttribute(deviceId, DataConstants.CLIENT_SCOPE, e.getKey()); if (latestOpt.isPresent()) { AttributeKvEntry latestEntry = latestOpt.get(); diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRpcMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRpcMsgHandler.java index b166dae43f..06467fe1a1 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRpcMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRpcMsgHandler.java @@ -97,7 +97,7 @@ public class TelemetryRpcMsgHandler implements RpcMsgHandler { builder.setDeviceId(cmd.getDeviceId().toString()); builder.setType(cmd.getType().name()); builder.setAllKeys(cmd.isAllKeys()); - cmd.getKeyStates().entrySet().stream().forEach(e -> builder.addKeyStates(SubscriptionKetStateProto.newBuilder().setKey(e.getKey()).setTs(e.getValue()).build())); + cmd.getKeyStates().entrySet().forEach(e -> builder.addKeyStates(SubscriptionKetStateProto.newBuilder().setKey(e.getKey()).setTs(e.getValue()).build())); ctx.sendPluginRpcMsg(new RpcMsg(address, SUBSCRIPTION_CLAZZ, builder.build().toByteArray())); } @@ -144,7 +144,7 @@ public class TelemetryRpcMsgHandler implements RpcMsgHandler { if (update.getErrorMsg() != null) { builder.setErrorMsg(update.getErrorMsg()); } - update.getData().entrySet().stream().forEach( + update.getData().entrySet().forEach( e -> { SubscriptionUpdateValueListProto.Builder dataBuilder = SubscriptionUpdateValueListProto.newBuilder(); @@ -166,7 +166,7 @@ public class TelemetryRpcMsgHandler implements RpcMsgHandler { return new SubscriptionUpdate(proto.getSubscriptionId(), SubscriptionErrorCode.forCode(proto.getErrorCode()), proto.getErrorMsg()); } else { Map> data = new TreeMap<>(); - proto.getDataList().stream().forEach(v -> { + proto.getDataList().forEach(v -> { List values = data.get(v.getKey()); if (values == null) { values = new ArrayList<>(); diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java index 6ea7489c5e..8e2d62abeb 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java @@ -109,8 +109,8 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { sendWsMsg(ctx, sessionRef, new SubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(keys.size()); - keys.stream().forEach(key -> subState.put(key, 0L)); - attributesData.stream().forEach(v -> subState.put(v.getKey(), v.getTs())); + keys.forEach(key -> subState.put(key, 0L)); + attributesData.forEach(v -> subState.put(v.getKey(), v.getTs())); sub = new SubscriptionState(sessionId, cmd.getCmdId(), deviceId, SubscriptionType.ATTRIBUTES, false, subState); } else { @@ -119,7 +119,7 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { sendWsMsg(ctx, sessionRef, new SubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(attributesData.size()); - attributesData.stream().forEach(v -> subState.put(v.getKey(), v.getTs())); + attributesData.forEach(v -> subState.put(v.getKey(), v.getTs())); sub = new SubscriptionState(sessionId, cmd.getCmdId(), deviceId, SubscriptionType.ATTRIBUTES, true, subState); } @@ -154,8 +154,8 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { sendWsMsg(ctx, sessionRef, new SubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(keys.size()); - keys.stream().forEach(key -> subState.put(key, startTs)); - data.stream().forEach(v -> subState.put(v.getKey(), v.getTs())); + keys.forEach(key -> subState.put(key, startTs)); + data.forEach(v -> subState.put(v.getKey(), v.getTs())); SubscriptionState sub = new SubscriptionState(sessionId, cmd.getCmdId(), deviceId, SubscriptionType.TIMESERIES, false, subState); subscriptionManager.addLocalWsSubscription(ctx, sessionId, deviceId, sub); } else { @@ -168,8 +168,8 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { sendWsMsg(ctx, sessionRef, new SubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(keys.size()); - keys.stream().forEach(key -> subState.put(key, startTs)); - data.stream().forEach(v -> subState.put(v.getKey(), v.getTs())); + keys.forEach(key -> subState.put(key, startTs)); + data.forEach(v -> subState.put(v.getKey(), v.getTs())); SubscriptionState sub = new SubscriptionState(sessionId, cmd.getCmdId(), deviceId, SubscriptionType.TIMESERIES, false, subState); subscriptionManager.addLocalWsSubscription(ctx, sessionId, deviceId, sub); } @@ -188,7 +188,7 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { public void onSuccess(PluginContext ctx, List data) { sendWsMsg(ctx, sessionRef, new SubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(data.size()); - data.stream().forEach(v -> subState.put(v.getKey(), v.getTs())); + data.forEach(v -> subState.put(v.getKey(), v.getTs())); SubscriptionState sub = new SubscriptionState(sessionId, cmd.getCmdId(), deviceId, SubscriptionType.TIMESERIES, true, subState); subscriptionManager.addLocalWsSubscription(ctx, sessionId, deviceId, sub); } diff --git a/extensions/extension-kafka/src/main/java/org/thingsboard/server/extensions/kafka/plugin/KafkaPlugin.java b/extensions/extension-kafka/src/main/java/org/thingsboard/server/extensions/kafka/plugin/KafkaPlugin.java index 1642fb5f3f..7321ad793f 100644 --- a/extensions/extension-kafka/src/main/java/org/thingsboard/server/extensions/kafka/plugin/KafkaPlugin.java +++ b/extensions/extension-kafka/src/main/java/org/thingsboard/server/extensions/kafka/plugin/KafkaPlugin.java @@ -47,7 +47,7 @@ public class KafkaPlugin extends AbstractPlugin { properties.put("buffer.memory", configuration.getBufferMemory()); if (configuration.getOtherProperties() != null) { configuration.getOtherProperties() - .stream().forEach(p -> properties.put(p.getKey(), p.getValue())); + .forEach(p -> properties.put(p.getKey(), p.getValue())); } init(); }