【问题标题】:How to recover from Exception thrown in Akka Streams Sink?如何从 Akka Streams Sink 中抛出的异常中恢复?
【发布时间】:2020-06-27 16:13:07
【问题描述】:

如何从 Akka Streams 接收器中抛出的异常中恢复?

简单示例:

    Source<Integer, NotUsed> integerSource = Source.from(Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9));

    integerSource.runWith(Sink.foreach(x -> {
      if (x == 4) {
        throw new Exception("Error Occurred");
      }
      System.out.println("Sink: " + x);
    }), system);

输出:

Sink: 1
Sink: 2
Sink: 3

如何处理异常并从源移至下一个元素? (又名 5,6,7,8,9)

【问题讨论】:

  • 您想跳过特定项目吗?然后只需filter他们(更多详情请参阅:doc.akka.io/docs/akka/current/stream/operators/Source-or-Flow/…)
  • 我不知道哪个元素会导致异常。我的 Sink 可能是 Kafka 或 RabbitMq,所以如果说与 rabbbitmq 的连接暂时失败,那么我不想停止流,而是继续未来的元素

标签: java scala akka akka-stream


【解决方案1】:

默认情况下,supervision strategy 在抛出异常时会停止流。要更改监督策略以删除导致异常的消息并继续处理下一条消息,请使用“恢复”策略。例如:

final Function<Throwable, Supervision.Directive> decider =
  exc -> {
    return Supervision.resume();
  };

final Sink<Integer, CompletionStage<Done>> printSink =
  Sink.foreach(x -> {
    if (x == 4) {
      throw new Exception("Error Occurred");
    }
    System.out.println("Sink: " + x);
  });

final RunnableGraph<CompletionStage<Done>> runnableGraph =
  integerSource.toMat(printSink, Keep.right());

final RunnableGraph<CompletionStage<Done>> withResumingSupervision =
  runnableGraph.withAttributes(ActorAttributes.withSupervisionStrategy(decider));

final CompletionStage<Done> result = withResumingSupervision.run(system);

您还可以为不同类型的异常定义不同的监督策略:

final Function<Throwable, Supervision.Directive> decider =
  exc -> {
    if (exc instanceof MySpecificException) return Supervision.resume();
    else return Supervision.stop();
  };

【讨论】:

  • 谢谢杰夫。它起作用了一个小的更正: withResumingSupervision 应该是 RunnableGraph> 类型
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-01-28
  • 1970-01-01
  • 2021-08-07
  • 1970-01-01
  • 2021-01-04
  • 2017-01-11
  • 1970-01-01
相关资源
最近更新 更多