diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java index ec7f491fed..632ce16bc1 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java @@ -320,28 +320,19 @@ public class RemoteTransportService extends AbstractTransportService { send(sessionInfo, toRuleEngineMsg, callback); } - private static class TransportCallbackAdaptor implements Callback { - private final TransportServiceCallback callback; - - TransportCallbackAdaptor(TransportServiceCallback callback) { - this.callback = callback; - } - - @Override - public void onCompletion(RecordMetadata metadata, Exception exception) { - if (exception == null) { - if (callback != null) { - callback.onSuccess(null); - } - } else { - if (callback != null) { - callback.onError(exception); + private void send(SessionInfoProto sessionInfo, ToRuleEngineMsg toRuleEngineMsg, TransportServiceCallback callback) { + ruleEngineProducer.send(getRoutingKey(sessionInfo), toRuleEngineMsg, (metadata, exception) -> { + if (callback != null) { + if (exception == null) { + this.transportCallbackExecutor.submit(() -> { + callback.onSuccess(null); + }); + } else { + this.transportCallbackExecutor.submit(() -> { + callback.onError(exception); + }); } } - } - } - - private void send(SessionInfoProto sessionInfo, ToRuleEngineMsg toRuleEngineMsg, TransportServiceCallback callback) { - ruleEngineProducer.send(getRoutingKey(sessionInfo), toRuleEngineMsg, new TransportCallbackAdaptor(callback)); + }); } }