如何用RxJava的Buffer收集AMQP事件并维持订阅持续执行?
Hey there, let's break down why your current code only prints empty lists and fix it properly.
The Root Problem
Your initial obs is an Observable.empty()—and when you subscribe to it, that's the exact stream your subscription is listening to. Every time you call obs = obs.concatWith(Observable.just(event)), you're creating a new Observable instance and updating the obs variable, but the original subscription (the one stored in disposable) is still tied to the empty stream. That's why you never see any events show up in your buffer—it's waiting on a stream that will never emit anything.
Observable instances are immutable once created; reassigning the variable doesn't change the stream the subscription is watching.
The Solution: Use a Subject
You need a way to dynamically push events into an existing stream that your subscription is already listening to. That's exactly what a PublishSubject is for—it acts as both an Observable (so you can subscribe to it) and an Observer (so you can push events into it).
Here's how to rewrite your code:
// Replace the empty Observable with a PublishSubject private final PublishSubject<Event> eventSubject = PublishSubject.create(); private final Disposable disposable; // Initialize the subscription in a constructor or init block public YourConsumerClass() { disposable = eventSubject .buffer(10, SECONDS) .retry(t -> true) .subscribe(System.out::println); } @Override public void handle(final Event event, final MessageContext context) throws MessageConsumptionException { // Push the incoming event directly into the subject eventSubject.onNext(event); }
What Changed?
- PublishSubject: This maintains a single stream that your subscription listens to. Any events you push via
onNext()will flow through to the buffer operator. - No more reassigning observables: We don't need to create new Observable instances every time an event comes in—we just send the event to the existing subject.
- Persistent subscription: The subscription is tied to the subject from the start, so it will receive all events pushed into it after subscription.
Bonus Notes
- Make sure to handle cleanup properly: when your consumer is shut down, call
disposable.dispose()andeventSubject.onComplete()to avoid memory leaks. - If you ever need to retain the last event for new subscribers (not necessary here), you could use a
BehaviorSubjectinstead—butPublishSubjectis the perfect fit for this use case since you just want to pass through events as they arrive.
内容的提问来源于stack exchange,提问作者lmarx

