You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在RxJava中创建Observable与Observer?服务端推送数据处理场景

Hey there! Let's walk through this thoroughly—first covering the core ways to create Observables and Observers in RxJava, then tackling how to wrap your server-listening listen function into the RxJava pattern so you can use onNext, onComplete, and onError to handle the data flow.


1. Creating Observables in RxJava

Observables are the source of events in RxJava. Here are the most common ways to create them:

  • Observable.create(): The most flexible option, giving you full control over event emission
Observable<String> customObservable = Observable.create(emitter -> {
    try {
        emitter.onNext("First server update");
        emitter.onNext("Second server update");
        emitter.onComplete(); // Signal that no more events will be sent
    } catch (Exception e) {
        emitter.onError(e); // Forward any errors to subscribers
    }
});
  • Observable.just(): Emit a fixed set of pre-known values
// Emits 1, 2, 3 in sequence
Observable<Integer> justObservable = Observable.just(1, 2, 3);
  • Observable.fromXXX(): Create from collections, arrays, or Iterables
List<String> messageQueue = Arrays.asList("Ping", "Status Update", "Disconnect");
Observable<String> fromListObservable = Observable.fromIterable(messageQueue);
  • Timed Observables: For periodic or delayed events (great for polling, though not your use case here)
// Emits a long value every 1 second
Observable<Long> intervalObservable = Observable.interval(1, TimeUnit.SECONDS);

2. Creating Observers in RxJava

Observers subscribe to Observables and handle the emitted events. You have two main approaches:

  • Full Observer Implementation: For when you need all lifecycle hooks
Observer<String> serverObserver = new Observer<String>() {
    private Disposable disposable;

    @Override
    public void onSubscribe(Disposable d) {
        // Store the disposable to cancel the subscription later if needed
        disposable = d;
    }

    @Override
    public void onNext(String data) {
        // Handle incoming server data
        System.out.println("Received from server: " + data);
    }

    @Override
    public void onError(Throwable e) {
        // Handle errors (e.g., server connection failure)
        System.err.println("Server error occurred: " + e.getMessage());
    }

    @Override
    public void onComplete() {
        // Triggered when the Observable finishes emitting events (e.g., server closed connection)
        System.out.println("Server stream completed");
    }
};
  • Simplified Subscription with Lambdas: Skip the full Observer implementation for simpler use cases
customObservable.subscribe(
    data -> System.out.println("Received: " + data), // onNext handler
    error -> System.err.println("Error: " + error.getMessage()), // onError handler
    () -> System.out.println("Stream completed") // onComplete handler
);

3. Wrapping Your listen Function into RxJava

Since your listen function is a void callback-based method, we'll wrap it using Observable.create() to bridge it into the RxJava ecosystem. First, let's assume your listen function uses a callback interface (this is standard for push-based APIs):

// Assume this is the callback interface your listen function uses
interface ServerPushCallback {
    void onDataReceived(String serverData);
    void onConnectionClosed();
    void onListenFailure(Throwable error);
}

// Your existing void listen function
void listen(ServerPushCallback callback) {
    // Internal logic: connects to server, triggers callbacks on data/error/close
}

Now, let's wrap this into an Observable that emits server data, errors, and completion signals:

Observable<String> serverPushObservable = Observable.create(emitter -> {
    ServerPushCallback callback = new ServerPushCallback() {
        @Override
        public void onDataReceived(String serverData) {
            // Only emit if the subscription is still active (avoids memory leaks)
            if (!emitter.isDisposed()) {
                emitter.onNext(serverData);
            }
        }

        @Override
        public void onConnectionClosed() {
            if (!emitter.isDisposed()) {
                emitter.onComplete();
            }
        }

        @Override
        public void onListenFailure(Throwable error) {
            if (!emitter.isDisposed()) {
                emitter.onError(error);
            }
        }
    };

    // Start listening to the server
    listen(callback);

    // Define cleanup logic for when the subscription is canceled
    emitter.setCancellable(() -> {
        // Add code here to stop listening to the server (e.g., disconnect)
        // stopListeningToServer();
    });
});

Finally, subscribe to this Observable to handle the server data:

// Using the full Observer
serverPushObservable.subscribe(serverObserver);

// Or using simplified lambdas
serverPushObservable.subscribe(
    serverData -> System.out.println("Server pushed: " + serverData),
    error -> System.err.println("Server error: " + error.getMessage()),
    () -> System.out.println("Server connection closed gracefully")
);

Key notes here:

  • Always check emitter.isDisposed() before emitting events to avoid unnecessary work or memory leaks.
  • Use emitter.setCancellable() to clean up resources (like disconnecting from the server) when the subscription is canceled.

内容的提问来源于stack exchange,提问作者jenjen

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 07:53:13