Spring Cache抽象与HystrixObservableCommand集成异常问题
Hey there, let's break down why you're hitting that HystrixRuntimeException and fix it up.
异常根源
The core issue here is that a single HystrixObservableCommand instance can only be executed once, but Spring Cache's @Cacheable annotation is causing multiple subscriptions to the same Observable returned by your getUserList method.
When Spring Cache processes an Observable return value, it may subscribe to it multiple times (e.g., first to check if it can cache the result, then again to serve the value to the caller). Since your current code creates one HystrixObservableCommand per getUserList call and returns its Observable directly, subsequent subscriptions trigger Hystrix's protection against reusing command instances.
Also, a quick note: your construct and resumeWithFallback methods don't call subscriber.onNext() to emit data—this would lead to empty results even if the exception was fixed, so we'll address that too.
Fixes to Try
1. Use Observable.defer() to Create New Hystrix Commands Per Subscription
This is the cleanest solution for keeping your reactive flow intact. By wrapping your Hystrix command creation in Observable.defer(), you ensure every subscription gets a brand new HystrixObservableCommand instance, avoiding the "multiple execution" violation.
Update your UserRepositoryImpl.getUserList() method like this:
@Override public Observable<Optional<List<User>>> getUserList(final String status) { // Defer command creation until subscription time return Observable.defer(() -> { return new HystrixObservableCommand<Optional<List<User>>>( HystrixObservableCommand.Setter.withGroupKey(hystrixCommandGroupKey) .andCommandKey(HystrixCommandKey.Factory.asKey("getUserList"))) { @Override protected Observable<Optional<List<User>>> construct() { return Observable.create((Subscriber<? super Optional<List<User>>> subscriber) -> { try { if (!subscriber.isUnsubscribed()) { System.out.println("##############Database Call getUserList"); // Replace with your actual DB/HTTP call logic Optional<List<User>> result = /* fetch data */; subscriber.onNext(result); // Emit the fetched data } } catch (Exception ex) { subscriber.onError(ex); // Properly propagate exceptions } finally { if (!subscriber.isUnsubscribed()) { subscriber.onCompleted(); } } }); } @Override protected Observable<Optional<List<User>>> resumeWithFallback() { return Observable.create((Subscriber<? super Optional<List<User>>> subscriber) -> { try { if (!subscriber.isUnsubscribed()) { System.out.println("##############Fallback for getUserList"); // Your fallback logic (e.g., return empty list) subscriber.onNext(Optional.empty()); } } catch (Exception ex) { subscriber.onError(ex); } finally { if (!subscriber.isUnsubscribed()) { subscriber.onCompleted(); } } }); } }.toObservable(); }); }
2. Convert Observable to CompletableFuture (For Non-Strict Reactive Scenarios)
If you don't need full reactive behavior, you can convert the Observable to a CompletableFuture before caching. This way, Spring Cache only interacts with the Future once, avoiding multiple subscriptions:
Update your UserService.listUsersAsync() method:
@Cacheable(cacheNames = "UserListCache", key = "#status") public CompletableFuture<Optional<List<User>>> listUsersAsync(final String status){ return userRepository.getUserList(status) .toBlocking() .toFuture(); }
Note: This blocks the Observable execution, so it's not ideal if you need to keep the reactive pipeline fully non-blocking.
3. Adjust Spring Cache Handling for Reactive Types
If you want to stick with Observables, you can also look into customizing how Spring Cache interacts with reactive streams. For example, using cache() operator on the Observable to make it replayable, but you still need to combine this with Observable.defer() to ensure fresh Hystrix commands per subscription.
Quick Additional Checks
- Make sure your
hystrixCommandGroupKeyis properly initialized (e.g., not a shared instance causing unintended command grouping). - In your fallback method, avoid repeating the same DB/HTTP call as the
constructmethod—fallbacks should be lightweight (like returning default data) instead of retrying the same operation.
内容的提问来源于stack exchange,提问作者Ankit8051

