guys!
I'm developing an android application and want to integrate notifications using stomp over websockets. At the backend we use Spring Boot app with Spring websockets & RabbitMQ and Stomp.
I'm using the specified library StompProtocolAndroid 1.6.6 and faced with issue below:
Source code:
public void connectStomp() {
mStompClient = Stomp.over(Stomp.ConnectionProvider.OKHTTP, ANDROID_EMULATOR_LOCALHOST);
mStompClient.withClientHeartbeat(1000).withServerHeartbeat(1000);
String destinationPath = "/exchange/push.tx/push-data.userId" + userId;
resetSubscriptions();
mStompClient.send("", "").subscribe();
Disposable dispLifecycle = mStompClient.lifecycle()
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(lifecycleEvent -> {
switch (lifecycleEvent.getType()) {
case OPENED:
break;
case ERROR:
Log.e(TAG, "Stomp connection error", lifecycleEvent.getException());
resetSubscriptions();
mStompClient.disconnect();
mStompClient.connect();
break;
case CLOSED:
resetSubscriptions();
break;
case FAILED_SERVER_HEARTBEAT:
break;
}
});
compositeDisposable.add(dispLifecycle);
// Receive greetings
Disposable dispTopic = mStompClient.topic(destinationPath)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(topicMessage -> {
Log.d(TAG, "Received " + topicMessage.getPayload());
addNotification(gson.fromJson(topicMessage.getPayload(), Notification.class));
}, throwable -> {
Log.e(TAG, "Error on subscribe topic", throwable);
});
compositeDisposable.add(dispTopic);
mStompClient.connect();
}
private void resetSubscriptions() {
if (compositeDisposable != null) {
compositeDisposable.dispose();
}
compositeDisposable = new CompositeDisposable();
}
public void disconnectStomp() {
mStompClient = Stomp.over(Stomp.ConnectionProvider.OKHTTP, ANDROID_EMULATOR_LOCALHOST);
mStompClient.disconnect();
}
Here is the stacktrace:
AbstractConnectionProvider: Receive STOMP message: ERROR
message:Connection to broker closed.
content-length:0
E/StompClient: Error parsing message
java.util.NoSuchElementException
at java.util.Scanner.skip(Scanner.java:1755)
at java.util.Scanner.skip(Scanner.java:1772)
at ua.naiksoftware.stomp.dto.StompMessage.from(StompMessage.java:89)
at ua.naiksoftware.stomp.StompClient$$ExternalSyntheticLambda3.apply(Unknown Source:2)
at io.reactivex.internal.operators.observable.ObservableMap$MapObserver.onNext(ObservableMap.java:57)
at io.reactivex.internal.operators.observable.ObservableConcatMap$ConcatMapDelayErrorObserver$DelayErrorInnerObserver.onNext(ObservableConcatMap.java:506)
at io.reactivex.subjects.PublishSubject$PublishDisposable.onNext(PublishSubject.java:308)
at io.reactivex.subjects.PublishSubject.onNext(PublishSubject.java:228)
at ua.naiksoftware.stomp.provider.AbstractConnectionProvider.emitMessage(AbstractConnectionProvider.java:117)
at ua.naiksoftware.stomp.provider.OkHttpConnectionProvider$1.onMessage(OkHttpConnectionProvider.java:67)
at okhttp3.internal.ws.RealWebSocket.onReadMessage(RealWebSocket.java:322)
at okhttp3.internal.ws.WebSocketReader.readMessageFrame(WebSocketReader.java:219)
at okhttp3.internal.ws.WebSocketReader.processNextFrame(WebSocketReader.java:105)
at okhttp3.internal.ws.RealWebSocket.loopReader(RealWebSocket.java:273)
at okhttp3.internal.ws.RealWebSocket$1.onResponse(RealWebSocket.java:209)
at okhttp3.RealCall$AsyncCall.execute(RealCall.java:174)
at okhttp3.internal.NamedRunnable.run(NamedRunnable.java:32)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1167)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:641)
at java.lang.Thread.run(Thread.java:919)
D/AbstractConnectionProvider: Emit lifecycle event: CLOSED
D/StompClient: Socket closed
D/StompClient: Stomp disconnected
D/StompClient: Unsubscribe path: /exchange/push.tx/push.receive.inspectorId-23232id: 43ca8f76-5980-445d-ba87-91f9e3b4a60a