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

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 hystrixCommandGroupKey is 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 construct method—fallbacks should be lightweight (like returning default data) instead of retrying the same operation.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 03:57:40