爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      【SpringCloud技术专题】「Hystrix」(3)Command运作的原理和源码分析

      首页 知识中心 软件开发 文章详情页

      【SpringCloud技术专题】「Hystrix」(3)Command运作的原理和源码分析

      2023-06-19 06:58:10 阅读次数:450

      SpringCloud,微服务

      [每日一句]

      也许你度过了很糟糕的一天,但这并不代表你会因此度过糟糕的一生。

      构建一个Hystrix的Command模式

      这里我们需要关注三点:

      • (模板构造器)HystrixCommand构造函数当中的super

      • (真正的执行者)HystrixCommand定义的run,run其实就是真正执行命令的地方

      • (触发启动)new HelloWorldHystrixCommand("test").execute()中execute是发起执行的过程

      实现Demo

      public class HelloWorldHystrixCommand extends HystrixCommand<String>{
      	private final String name;
          public HelloWorldHystrixCommand(String name) {
              super(HystrixCommandGroupKey.Factory.asKey("ExampleGroup"));
              this.name = name;
          }
       
          @Override
          protected String run() throws Exception {
              //Thread.sleep(100);
              return "hello"+name;
          }
      }
      
      public static void main(String[] args){
          String result = new HelloWorldHystrixCommand("test").execute();
          System.out.println(result);
      }
      

      HystrixCommand初始化过程

      HystrixCommand的类关系图下,这里我们只需要暂时关注HystrixCommand继承自AbstractCommand即可,其他的我也没仔细看。

      HystrixCommand类依赖图

      【SpringCloud技术专题】「Hystrix」(3)Command运作的原理和源码分析

      HelloWorldHystrixCommand的构造步骤如下:
      1. 具体类HelloWorldHystrixCommand继承自HystrixCommand, 通过super()调用了HystrixCommand的构造函数
      2. HystrixCommand通过super()命令调用AbstractCommand实现初始化

      AbstractCommand类当中比较核心的几个对象如下:

      • metrics:统计指标
      • circuitBreaker:熔断器变量
      • threadPool:隔离的线程池
      • concurrencyStrategy :并发策略
      protected HystrixCommand(HystrixCommandGroupKey group) {
              super(group, null, null, null, null, null, null, null, null, null, null, null);
       }
      
      protected AbstractCommand(HystrixCommandGroupKey group, HystrixCommandKey key, 
                  HystrixThreadPoolKey threadPoolKey, HystrixCircuitBreaker circuitBreaker,
                  HystrixThreadPool threadPool,HystrixCommandProperties.Setter commandPropertiesDefaults, 
                  HystrixThreadPoolProperties.Setter threadPoolPropertiesDefaults,
                  HystrixCommandMetrics metrics, TryableSemaphore fallbackSemaphore,
                  TryableSemaphore executionSemaphore,
                  HystrixPropertiesStrategy propertiesStrategy,
      		    HystrixCommandExecutionHook executionHook) {
              this.commandGroup = initGroupKey(group);
              this.commandKey = initCommandKey(key, getClass());
              this.properties = initCommandProperties(this.commandKey, propertiesStrategy, commandPropertiesDefaults);
              this.threadPoolKey = initThreadPoolKey(threadPoolKey, this.commandGroup, 
                                                this.properties.executionIsolationThreadPoolKeyOverride().get());
              this.metrics = initMetrics(metrics, this.commandGroup, this.threadPoolKey, this.commandKey, this.properties);
              this.circuitBreaker = initCircuitBreaker(this.properties.circuitBreakerEnabled().get(), 
                                  circuitBreaker, this.commandGroup, this.commandKey, this.properties, this.metrics);
      
              //  线程池相关配置,通过线程池进行隔离
              this.threadPool = initThreadPool(threadPool, this.threadPoolKey, threadPoolPropertiesDefaults);
      
              //Strategies from plugins
              this.eventNotifier = HystrixPlugins.getInstance().getEventNotifier();
              this.concurrencyStrategy = HystrixPlugins.getInstance().getConcurrencyStrategy();
              HystrixMetricsPublisherFactory.createOrRetrievePublisherForCommand(this.commandKey, this.commandGroup, this.metrics, this.circuitBreaker, this.properties);
              this.executionHook = initExecutionHook(executionHook);
      
              this.requestCache = HystrixRequestCache.getInstance(this.commandKey, this.concurrencyStrategy);
              this.currentRequestLog = initRequestLog(this.properties.requestLogEnabled().get(), this.concurrencyStrategy);
      
              /* fallback semaphore override if applicable */
              this.fallbackSemaphoreOverride = fallbackSemaphore;
      
              /* execution semaphore override if applicable */
              this.executionSemaphoreOverride = executionSemaphore;
      }
      

      Hystrix的执行

      从Hystrix的整个执行的生命周期来看,可以分为两个阶段,阶段一主要是Observable的创建,阶段二主要是Observable的执行。

      • 两个过程的实际实现中运用了大量的RxJava的技能包,所以阅读起来有一点绕,我只能按照我粗浅的理解来尽量把整个过程讲解清楚。

      • 切入详细的过程当中,大家需要带有两个疑问去看代码,只有找到能解答这两个疑问的代码才算看懂了主流程,两个疑问分别是:(1)如何分配执行线程;(2)如何判定超时。

      Hystrix的Observable创建过程

      Hystrix的创建过程比较复杂,大致核心流程如下:

      1. 类HystrixCommand中execute方法开始执行,内部的queue()是实际执行整个过程,get()是获取执行的结果。

      2. 类HystrixCommand中queue方法,delegate = toObservable().toBlocking().toFuture(),toObservable负责创建Observable对象,toFuture负责执行任务(Future)。

      3. 类AbstractCommand中toObservable方法,hystrixObservable = Observable.defer(applyHystrixSemantics)负责关联applyHystrixSemantics。

      4. 类AbstractCommand中applyHystrixSemantics方法,executeCommandAndObserve(cmd)负责执行具体的AbstractCommand采用相关的Observable进行关联绑定。

      5. 类AbstractCommand中executeCommandAndObserve方法中,executeCommandWithSpecifiedIsolation(cmd).lift(new HystrixObservableTimeoutOperator<R>(_cmd))负责关联执行_cmd并关联超时检测任务。

      6. 类AbstractCommand中executeCommandWithSpecifiedIsolation是执行的具体的命令,【HystrixObservableTimeoutOperator是超时检测任务】。

      7. 类AbstractCommand中executeCommandWithSpecifiedIsolation方法中,getUserExecutionObservable负责执行具体任务,同时通过subscribeOn(threadPool.getScheduler(new Func0<Boolean>()))关联threadPool隔离执行任务,关键的隔离任务的位置。

      8. 类HystrixCommand中getUserExecutionObservable方法中,Observable.just(run())负责执行任务,这个run方法就是HelloWorldHystrixCommand的run方法,也就是这里终于回调回了真正的run函数。

      execute执行方法

      public R execute() {
              try {
                  return queue().get();
              } catch (Exception e) {
                  throw Exceptions.sneakyThrow(decomposeException(e));
              }
          }
      

      获取Observable获取相关的Future

      
      public Future<R> queue() {
              //todo ((Observable<T>)that).single().subscribe(new Subscriber<T>()
              //todo BlockingOperatorToFuture.toFuture里面真正执行任务
              //todo toObservable内部是通过RxJava构建
              final Future<R> delegate = toObservable().toBlocking().toFuture();
              final Future<R> f = new Future<R>() {
                  @Override
                  public boolean cancel(boolean mayInterruptIfRunning) {
                      if (delegate.isCancelled()) {
                          return false;
                      }
                      if (HystrixCommand.this.getProperties().executionIsolationThreadInterruptOnFutureCancel().get()) {
                          interruptOnFutureCancel.compareAndSet(false, mayInterruptIfRunning);
                      }
                      final boolean res = delegate.cancel(interruptOnFutureCancel.get());
      
                      if (!isExecutionComplete() && interruptOnFutureCancel.get()) {
                          final Thread t = executionThread.get();
                          if (t != null && !t.equals(Thread.currentThread())) {
                              t.interrupt();
                          }
                      }
      
                      return res;
                  }
      
                  @Override
                  public boolean isCancelled() {
                      return delegate.isCancelled();
                  }
      
                  @Override
                  public boolean isDone() {
                      return delegate.isDone();
                  }
      
                  @Override
                  public R get() throws InterruptedException, ExecutionException {
                      return delegate.get();
                  }
      
                  @Override
                  public R get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
                      return delegate.get(timeout, unit);
                  }
      
              };
      
              if (f.isDone()) {
                  try {
                      f.get();
                      return f;
                  } catch (Exception e) {
                      Throwable t = decomposeException(e);
                      if (t instanceof HystrixBadRequestException) {
                          return f;
                      } else if (t instanceof HystrixRuntimeException) {
                          HystrixRuntimeException hre = (HystrixRuntimeException) t;
                          switch (hre.getFailureType()) {
                          case COMMAND_EXCEPTION:
                          case TIMEOUT:
                              return f;
                          default:
                              throw hre;
                          }
                      } else {
                          throw Exceptions.sneakyThrow(t);
                      }
                  }
              }
      
              return f;
          }
      
      public Observable<R> toObservable() {
              final AbstractCommand<R> _cmd = this;
      
              final Action0 terminateCommandCleanup = new Action0() {
                  @Override
                  public void call() {
                      if (_cmd.commandState.compareAndSet(CommandState.OBSERVABLE_CHAIN_CREATED, CommandState.TERMINAL)) {
                          handleCommandEnd(false); //user code never ran
                      } else if (_cmd.commandState.compareAndSet(CommandState.USER_CODE_EXECUTED, CommandState.TERMINAL)) {
                          handleCommandEnd(true); //user code did run
                      }
                  }
              };
      
              final Action0 unsubscribeCommandCleanup = new Action0() {
                  @Override
                  public void call() {
                      circuitBreaker.markNonSuccess();
                      if (_cmd.commandState.compareAndSet(CommandState.OBSERVABLE_CHAIN_CREATED, CommandState.UNSUBSCRIBED)) {
                          if (!_cmd.executionResult.containsTerminalEvent()) {
                              _cmd.eventNotifier.markEvent(HystrixEventType.CANCELLED, _cmd.commandKey);
                              try {
                                  executionHook.onUnsubscribe(_cmd);
                              } catch (Throwable hookEx) {
                                  logger.warn("Error calling HystrixCommandExecutionHook.onUnsubscribe", hookEx);
                              }
                              _cmd.executionResultAtTimeOfCancellation = _cmd.executionResult
                                      .addEvent((int) (System.currentTimeMillis() - _cmd.commandStartTimestamp), 
                                      HystrixEventType.CANCELLED);
                          }
                          handleCommandEnd(false); //user code never ran
                      } else if (_cmd.commandState.compareAndSet(CommandState.USER_CODE_EXECUTED, CommandState.UNSUBSCRIBED)) {
                          if (!_cmd.executionResult.containsTerminalEvent()) {
                              _cmd.eventNotifier.markEvent(HystrixEventType.CANCELLED, _cmd.commandKey);
                              try {
                                  executionHook.onUnsubscribe(_cmd);
                              } catch (Throwable hookEx) {
                                  logger.warn("Error calling HystrixCommandExecutionHook.onUnsubscribe", hookEx);
                              }
                              _cmd.executionResultAtTimeOfCancellation = _cmd.executionResult
                                      .addEvent((int) (System.currentTimeMillis() - _cmd.commandStartTimestamp), 
                                      HystrixEventType.CANCELLED);
                          }
                          handleCommandEnd(true); //user code did run
                      }
                  }
              };
      
              final Func0<Observable<R>> applyHystrixSemantics = new Func0<Observable<R>>() {
                  @Override
                  public Observable<R> call() {
                      if (commandState.get().equals(CommandState.UNSUBSCRIBED)) {
                          return Observable.never();
                      }
                      return applyHystrixSemantics(_cmd);
                  }
              };
      
              final Func1<R, R> wrapWithAllOnNextHooks = new Func1<R, R>() {
                  @Override
                  public R call(R r) {
                      R afterFirstApplication = r;
      
                      try {
                          afterFirstApplication = executionHook.onComplete(_cmd, r);
                      } catch (Throwable hookEx) {
                          logger.warn("Error calling HystrixCommandExecutionHook.onComplete", hookEx);
                      }
      
                      try {
                          return executionHook.onEmit(_cmd, afterFirstApplication);
                      } catch (Throwable hookEx) {
                          logger.warn("Error calling HystrixCommandExecutionHook.onEmit", hookEx);
                          return afterFirstApplication;
                      }
                  }
              };
      
              final Action0 fireOnCompletedHook = new Action0() {
                  @Override
                  public void call() {
                      try {
                          executionHook.onSuccess(_cmd);
                      } catch (Throwable hookEx) {
                          logger.warn("Error calling HystrixCommandExecutionHook.onSuccess", hookEx);
                      }
                  }
              };
      
              //defer在subscribe的时候会真正执行
              return Observable.defer(new Func0<Observable<R>>() {
                  @Override
                  public Observable<R> call() {
                       /* this is a stateful object so can only be used once */
                      if (!commandState.compareAndSet(CommandState.NOT_STARTED, CommandState.OBSERVABLE_CHAIN_CREATED)) {
                          IllegalStateException ex = new IllegalStateException("");
                          //TODO make a new error type for this
                          throw new HystrixRuntimeException(FailureType.BAD_REQUEST_EXCEPTION, 
                             _cmd.getClass(), getLogMessagePrefix() + " ", ex, null);
                      }
      
                      commandStartTimestamp = System.currentTimeMillis();
      
                      if (properties.requestLogEnabled().get()) {
                          if (currentRequestLog != null) {
                              currentRequestLog.addExecutedCommand(_cmd);
                          }
                      }
      
                      final boolean requestCacheEnabled = isRequestCachingEnabled();
                      final String cacheKey = getCacheKey();
      
                      /* try from cache first */
                      if (requestCacheEnabled) {
                          HystrixCommandResponseFromCache<R> fromCache = 
      (HystrixCommandResponseFromCache<R>) requestCache.get(cacheKey);
                          if (fromCache != null) {
                              isResponseFromCache = true;
                              return handleRequestCacheHitAndEmitValues(fromCache, _cmd);
                          }
                      }
      
                      // todo 这里应该在会去执行applyHystrixSemantics方法
                      Observable<R> hystrixObservable =
                              Observable.defer(applyHystrixSemantics)
                                      .map(wrapWithAllOnNextHooks);
      
                      Observable<R> afterCache;
      
                      // put in cache
                      if (requestCacheEnabled && cacheKey != null) {
                          // wrap it for caching
                          HystrixCachedObservable<R> toCache = HystrixCachedObservable.from(hystrixObservable, _cmd);
                          HystrixCommandResponseFromCache<R> fromCache = 
      (HystrixCommandResponseFromCache<R>) requestCache.putIfAbsent(cacheKey, toCache);
                          if (fromCache != null) {
                              // another thread beat us so we'll use the cached value instead
                              toCache.unsubscribe();
                              isResponseFromCache = true;
                              return handleRequestCacheHitAndEmitValues(fromCache, _cmd);
                          } else {
                              // we just created an ObservableCommand so we cast and return it
                              afterCache = toCache.toObservable();
                          }
                      } else {
                          afterCache = hystrixObservable;
                      }
      
                      //todo 关联逻辑
                      return afterCache
                              .doOnTerminate(terminateCommandCleanup)    
                              .doOnUnsubscribe(unsubscribeCommandCleanup) 
                              .doOnCompleted(fireOnCompletedHook);
                  }
              });
          }
      
          // TODO: 2018/7/9 真正执行代码的地方在这里
          private Observable<R> applyHystrixSemantics(final AbstractCommand<R> _cmd) {
              executionHook.onStart(_cmd);
      
              if (circuitBreaker.attemptExecution()) {
                  final TryableSemaphore executionSemaphore = getExecutionSemaphore();
                  final AtomicBoolean semaphoreHasBeenReleased = new AtomicBoolean(false);
                  final Action0 singleSemaphoreRelease = new Action0() {
                      @Override
                      public void call() {
                          if (semaphoreHasBeenReleased.compareAndSet(false, true)) {
                              executionSemaphore.release();
                          }
                      }
                  };
      
                  final Action1<Throwable> markExceptionThrown = new Action1<Throwable>() {
                      @Override
                      public void call(Throwable t) {
                          eventNotifier.markEvent(HystrixEventType.EXCEPTION_THROWN, commandKey);
                      }
                  };
      
                  if (executionSemaphore.tryAcquire()) {
                      try {
                          executionResult = executionResult.setInvocationStartTime(System.currentTimeMillis());
      
                          // TODO: 2018/7/9 executeCommandAndObserve执行任务的地方
                          return executeCommandAndObserve(_cmd)
                                  .doOnError(markExceptionThrown)
                                  .doOnTerminate(singleSemaphoreRelease)
                                  .doOnUnsubscribe(singleSemaphoreRelease);
                      } catch (RuntimeException e) {
                          return Observable.error(e);
                      }
                  } else {
                      return handleSemaphoreRejectionViaFallback();
                  }
              } else {
                  return handleShortCircuitViaFallback();
              }
          }
      
      private Observable<R> executeCommandAndObserve(final AbstractCommand<R> _cmd) {
              final HystrixRequestContext currentRequestContext = HystrixRequestContext.getContextForCurrentThread();
      
              final Action1<R> markEmits = new Action1<R>() {
                  @Override
                  public void call(R r) {
                      if (shouldOutputOnNextEvents()) {
                          executionResult = executionResult.addEvent(HystrixEventType.EMIT);
                          eventNotifier.markEvent(HystrixEventType.EMIT, commandKey);
                      }
                      if (commandIsScalar()) {
                          long latency = System.currentTimeMillis() - executionResult.getStartTimestamp();
                          eventNotifier.markEvent(HystrixEventType.SUCCESS, commandKey);
                          executionResult = executionResult.addEvent((int) latency, HystrixEventType.SUCCESS);
                          eventNotifier.markCommandExecution(getCommandKey(), 
      properties.executionIsolationStrategy().get(), (int) latency, executionResult.getOrderedList());
                          circuitBreaker.markSuccess();
                      }
                  }
              };
      
              final Action0 markOnCompleted = new Action0() {
                  @Override
                  public void call() {
                      if (!commandIsScalar()) {
                          long latency = System.currentTimeMillis() - executionResult.getStartTimestamp();
                          eventNotifier.markEvent(HystrixEventType.SUCCESS, commandKey);
                          executionResult = executionResult.addEvent((int) latency, HystrixEventType.SUCCESS);
                          eventNotifier.markCommandExecution(getCommandKey(), properties.executionIsolationStrategy().get(), 
      (int) latency, executionResult.getOrderedList());
                          circuitBreaker.markSuccess();
                      }
                  }
              };
      
              final Func1<Throwable, Observable<R>> handleFallback = new Func1<Throwable, Observable<R>>() {
                  @Override
                  public Observable<R> call(Throwable t) {
                      circuitBreaker.markNonSuccess();
                      Exception e = getExceptionFromThrowable(t);
                      executionResult = executionResult.setExecutionException(e);
                      if (e instanceof RejectedExecutionException) {
                          return handleThreadPoolRejectionViaFallback(e);
                      } else if (t instanceof HystrixTimeoutException) {
                          return handleTimeoutViaFallback();
                      } else if (t instanceof HystrixBadRequestException) {
                          return handleBadRequestByEmittingError(e);
                      } else {
                          /*
                           * Treat HystrixBadRequestException from ExecutionHook like a plain HystrixBadRequestException.
                           */
                          if (e instanceof HystrixBadRequestException) {
                              eventNotifier.markEvent(HystrixEventType.BAD_REQUEST, commandKey);
                              return Observable.error(e);
                          }
      
                          return handleFailureViaFallback(e);
                      }
                  }
              };
      
              final Action1<Notification<? super R>> setRequestContext = new Action1<Notification<? super R>>() {
                  @Override
                  public void call(Notification<? super R> rNotification) {
                      setRequestContextIfNeeded(currentRequestContext);
                  }
              };
      
              //todo 这里是否执行超时动作,executeCommandWithSpecifiedIsolation这个函数非常重要
              Observable<R> execution;
              if (properties.executionTimeoutEnabled().get()) {
                  execution = executeCommandWithSpecifiedIsolation(_cmd)
                          .lift(new HystrixObservableTimeoutOperator<R>(_cmd));
                  //todo HystrixObservableTimeoutOperator负责执行超时动作
              } else {
                  execution = executeCommandWithSpecifiedIsolation(_cmd);
              }
      
              return execution.doOnNext(markEmits)
                      .doOnCompleted(markOnCompleted)
                      .onErrorResumeNext(handleFallback)
                      .doOnEach(setRequestContext);
          }
      
      private Observable<R> executeCommandWithSpecifiedIsolation(final AbstractCommand<R> _cmd) {
              if (properties.executionIsolationStrategy().get() == ExecutionIsolationStrategy.THREAD) {
                  return Observable.defer(new Func0<Observable<R>>() {
                      @Override
                      public Observable<R> call() {
                          executionResult = executionResult.setExecutionOccurred();
                          if (!commandState.compareAndSet(CommandState.OBSERVABLE_CHAIN_CREATED, CommandState.USER_CODE_EXECUTED)) {
                              return Observable.error(new IllegalStateException(
                                       "execution attempted while in state : " + commandState.get().name()));
                          }
      
                          metrics.markCommandStart(commandKey, threadPoolKey, ExecutionIsolationStrategy.THREAD);
      
                          if (isCommandTimedOut.get() == TimedOutStatus.TIMED_OUT) {
                              // the command timed out in the wrapping thread so we will return immediately
                              // and not increment any of the counters below or other such logic
                              return Observable.error(new RuntimeException("timed out before executing run()"));
                          }
                          if (threadState.compareAndSet(ThreadState.NOT_USING_THREAD, ThreadState.STARTED)) {
                              //we have not been unsubscribed, so should proceed
                              HystrixCounters.incrementGlobalConcurrentThreads();
                              threadPool.markThreadExecution();
                              // store the command that is being run
                              endCurrentThreadExecutingCommand = Hystrix.startCurrentThreadExecutingCommand(getCommandKey());
                              executionResult = executionResult.setExecutedInThread();
      
                              try {
                                  //todo 估计是前置依赖吧
                                  executionHook.onThreadStart(_cmd);
                                  executionHook.onRunStart(_cmd);
                                  executionHook.onExecutionStart(_cmd);
      
                                  //todo 真正执行的地方
                                  return getUserExecutionObservable(_cmd);
                              } catch (Throwable ex) {
                                  return Observable.error(ex);
                              }
                          } else {
                              //command has already been unsubscribed, so return immediately
                              return Observable.empty();
                          }
                      }
                  }).doOnTerminate(new Action0() {
                      @Override
                      public void call() {
                          if (threadState.compareAndSet(ThreadState.STARTED, ThreadState.TERMINAL)) {
                              handleThreadEnd(_cmd);
                          }
                          if (threadState.compareAndSet(ThreadState.NOT_USING_THREAD, ThreadState.TERMINAL)) {
                          }
                      }
                  }).doOnUnsubscribe(new Action0() {
                      @Override
                      public void call() {
                          if (threadState.compareAndSet(ThreadState.STARTED, ThreadState.UNSUBSCRIBED)) {
                              handleThreadEnd(_cmd);
                          }
                          if (threadState.compareAndSet(ThreadState.NOT_USING_THREAD, ThreadState.UNSUBSCRIBED)) {
                              //if it was never started and was cancelled, then no need to clean up
                          }
                          //if it was terminal, then other cleanup handled it
                      }
                  }).subscribeOn(threadPool.getScheduler(new Func0<Boolean>() {
                      //todo subscribeOn 据说获取线程的地方????
                      @Override
                      public Boolean call() {
                          return properties.executionIsolationThreadInterruptOnTimeout().get() && 
                                    _cmd.isCommandTimedOut.get() == TimedOutStatus.TIMED_OUT;
                      }
                  }));
              } else {
                  return Observable.defer(new Func0<Observable<R>>() {
                      @Override
                      public Observable<R> call() {
                          executionResult = executionResult.setExecutionOccurred();
                          if (!commandState.compareAndSet(CommandState.OBSERVABLE_CHAIN_CREATED, CommandState.USER_CODE_EXECUTED)) {
                              return Observable.error(new IllegalStateException(
                                        "execution attempted while in state : " + commandState.get().name()));
                          }
      
                          metrics.markCommandStart(commandKey, threadPoolKey, ExecutionIsolationStrategy.SEMAPHORE);
                          // semaphore isolated
                          // store the command that is being run
                          endCurrentThreadExecutingCommand = Hystrix.startCurrentThreadExecutingCommand(getCommandKey());
                          try {
                              executionHook.onRunStart(_cmd);
                              executionHook.onExecutionStart(_cmd);
                              return getUserExecutionObservable(_cmd);  
                          } catch (Throwable ex) {
                              return Observable.error(ex);
                          }
                      }
                  });
              }
          }
      
      private Observable<R> getUserExecutionObservable(final AbstractCommand<R> _cmd) {
              Observable<R> userObservable;
              try {
                  userObservable = getExecutionObservable();
              } catch (Throwable ex) {
                  // the run() method is a user provided implementation so can throw instead of using Observable.onError
                  // so we catch it here and turn it into Observable.error
                  userObservable = Observable.error(ex);
              }        return userObservable
                      .lift(new ExecutionHookApplication(_cmd))
                      .lift(new DeprecatedOnRunHookApplication(_cmd));
          }
      
       // 真正执行run的位置
      final protected Observable<R> getExecutionObservable() {
              return Observable.defer(new Func0<Observable<R>>() {
                  @Override
                  public Observable<R> call() {
                      try {
                          return Observable.just(run());
                      } catch (Throwable ex) {
                          return Observable.error(ex);
                      }
                  }
              }).doOnSubscribe(new Action0() {
                  @Override
                  public void call() {
                      // Save thread on which we get subscribed so that we can interrupt it later if needed
                      executionThread.set(Thread.currentThread());
                  }
              });
          }
      
      

      Hystrix的执行过程

      toFuture过程中真正触发构建的Observable的的代码在((Observable<T>)that).single().subscribe()当中,关注几个方法:

      onCompleted负责设置完成标记。 onNext负责设置结果。

      public static <T> Future<T> toFuture(Observable<? extends T> that) {
      
        final CountDownLatch finished = new CountDownLatch(1);
              final AtomicReference<T> value = new AtomicReference<T>();
              final AtomicReference<Throwable> error = new AtomicReference<Throwable>();
              @SuppressWarnings("unchecked")
              final Subscription s = ((Observable<T>)that).single().subscribe(new Subscriber<T>(){
                  @Override
                  public void onCompleted() {
                      finished.countDown();
                  }
      
                  @Override
                  public void onError(Throwable e) {
                      error.compareAndSet(null, e);
                      finished.countDown();
                  }
      
                  @Override
                  public void onNext(T v) {
                  // "single" guarantees there is only one "onNext"
                      value.set(v);
                  }
              });
      
              return new Future<T>() {
                  private volatile boolean cancelled;
                  @Override
                  public boolean cancel(boolean mayInterruptIfRunning) {
                      if (finished.getCount() > 0) {
                          cancelled = true;
                          s.unsubscribe();
                          // release the latch (a race condition may have already released it by now)
                          finished.countDown();
                          return true;
                      } else {
                          // can't cancel
                          return false;
                      }
                  }
      
                  @Override
                  public boolean isCancelled() {
                      return cancelled;
                  }
      
                  @Override
                  public boolean isDone() {
                      return finished.getCount() == 0;
                  }
      
                  @Override
                  public T get() throws InterruptedException, ExecutionException {
                      finished.await();
                      return getValue();
                  }
      
                  @Override
                  public T get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
                      if (finished.await(timeout, unit)) {
                          return getValue();
                      } else {
                          throw new TimeoutException("Timed out after " + unit.toMillis(timeout) 
                                                                        + "ms waiting for underlying Observable.");
                      }
                  }
      
                  private T getValue() throws ExecutionException {
                      final Throwable throwable = error.get();
                      if (throwable != null) {
                          throw new ExecutionException("Observable onError", throwable);
                      } else if (cancelled) {
                          // Contract of Future.get() requires us to throw this:
                          throw new CancellationException("Subscription unsubscribed");
                      } else {
                          return value.get();
                      }
                  }
      
              };
      
          }
      
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.51cto.com/alex4dream/2943194,作者:洛神灬殇,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:SpringBoot-技术专区-整合mybatis 如何使用多数据源

      下一篇:设置SpringMVC默认首页启动

      相关文章

      2025-02-21 08:57:32

      【分布式】什么是服务熔断?什么是服务降级?

      【分布式】什么是服务熔断?什么是服务降级?

      2025-02-21 08:57:32
      服务 , 状态 , 调用 , 阈值
      2024-12-04 07:18:18

      SpringCloud-MQ消息队列

      MQ (MessageQueue) ,中文是消息队列,字面来看就是存放消息的队列。也就是事件驱动架构中的Broker。消息队列是一种基于生产者-消费者模型的通信方式,通过在消息队列中存放和传递消息,实现了不同组件、服务或系统之间的异步通信。

      2024-12-04 07:18:18
      SpringCloud , 队列
      2024-09-25 10:15:01

      十六、微服务之-耦合、耦合模理论、耦合谐振子模型、嵌入式工程实现低耦合的实例

      组件之间依赖关系强度的度量被称为耦合。好的设计总是高内聚和低耦合的。

      2024-09-25 10:15:01
      微服务
      2024-09-25 10:14:48

      Spring Cloud Alibaba入门五:openFeign实现REST调用-构造多参数请求

      Spring Cloud Alibaba入门五:openFeign实现REST调用-构造多参数请求

      2024-09-25 10:14:48
      class , http , SpringCloud
      2024-09-25 10:14:48

      Spring Cloud中使用LoggingClient来发送到LoggingAdmin记录日志

      在Spring Cloud部署方式下使用LoggingClient自动发现LoggingAdmin服务并上报日志,可以按照以下步骤进行操作。

      2024-09-25 10:14:48
      SpringCloud
      2024-09-25 10:14:48

      二十三、微服务之-【Spring Cloud中使用LoggingClient来发送到LoggingAdmin记录日志】

      在Spring Cloud部署方式下使用LoggingClient自动发现LoggingAdmin服务并上报日。

      2024-09-25 10:14:48
      微服务
      2024-09-25 10:14:21

      RabbitMq 消息确认机制详解 SpringCloud

      RabbitMq 消息确认机制详解 SpringCloud

      2024-09-25 10:14:21
      SpringCloud
      2024-09-25 10:14:09

      SpringCloud-技术专区-Gateway优雅的处理Filter抛出的异常

      Filter的位置相对比较尴尬,在MVC层之外,所以无法使用SpringMVC统一异常处理。     

      2024-09-25 10:14:09
      Spring , SpringCloud
      2024-09-24 06:29:45

      十七、微服务之-REST/RESTful

      Representational State Transfer(REST)/ RESTful (表述性状态转移)是一种帮助计算机系统通过 Internet 进行通信的架构风格。这使得微服务更容易理解和实现。

      2024-09-24 06:29:45
      Java , RESTful , 微服务
      2024-09-24 06:29:45

      微服务之-REST/RESTful

      Representational State Transfer(REST)/ RESTful (表述性状态转移)是一种帮助计算机系统通过 Internet 进行通信的架构风格。这使得微服务更容易理解和实现。

      2024-09-24 06:29:45
      RESTful , 微服务
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5220512

      查看更多

      最新文章

      SpringCloud-MQ消息队列

      2024-12-04 07:18:18

      Spring Cloud Alibaba入门五:openFeign实现REST调用-构造多参数请求

      2024-09-25 10:14:48

      Spring Cloud中使用LoggingClient来发送到LoggingAdmin记录日志

      2024-09-25 10:14:48

      二十三、微服务之-【Spring Cloud中使用LoggingClient来发送到LoggingAdmin记录日志】

      2024-09-25 10:14:48

      RabbitMq 消息确认机制详解 SpringCloud

      2024-09-25 10:14:21

      SpringCloud-技术专区-Gateway优雅的处理Filter抛出的异常

      2024-09-25 10:14:09

      查看更多

      热门文章

      浅谈SpringCloud 和 Dubbo 的区别

      2023-04-06 09:56:07

      RabbitMq 消息确认机制详解 SpringCloud

      2024-09-25 10:14:21

      SpringCloud ---- Eureka常见面试题

      2023-04-06 06:11:29

      k8s资源之APIService&ControllerRevision

      2023-02-16 09:40:38

      SpringCloud----Zuul-过滤器详解

      2023-04-06 06:11:29

      Spring项目&Spring-boot-starter 组件&Spring Cloud 生态圈以及微服务简介

      2023-06-16 06:12:13

      查看更多

      热门标签

      java Java python 编程开发 代码 开发语言 算法 线程 Python html 数组 C++ 元素 javascript c++
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      SpringCloud ---- Eureka常见面试题

      SpringCloud-技术专区-Gateway优雅的处理Filter抛出的异常

      SpringCloud 和springBoot基础注解及配置

      浅谈SpringCloud 和 Dubbo 的区别

      四十三、微服务之测试金字塔

      Spring Cloud中使用LoggingClient来发送到LoggingAdmin记录日志

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号