【发布时间】:2019-08-02 03:36:38
【问题描述】:
我们定义了一个基本订阅者,它通过抛出异常并依靠 Akka Streams 的流监督来恢复 Flow,从而跳过失败的消息(即出于某些业务逻辑原因,我们不会处理):
someLagomService
.someTopic()
.subscribe
.withGroupId("lagom-service")
.atLeastOnce(
Flow[Int]
.mapAsync(1)(el => {
// Exception may occur here or can map to Done
})
.withAttributes(ActorAttributes.supervisionStrategy({
case t =>
Supervision.Resume
})
)
这对于负载很小的基本用例来说似乎工作得很好,但我们注意到对于大量消息来说非常奇怪的事情(例如:非常频繁地重新处理消息等)。
深入研究代码,我们看到 Lagom 的 broker.Subscriber.atLeastOnce 文档指出:
flow可能会从上游提取更多元素,但它必须发出 对于它收到的每条消息,恰好有一条Done消息。它必须 也以与接收消息相同的顺序发出它们。这 意味着flow不得过滤或收集 消息,相反,它必须将消息拆分为单独的流,并 将那些将被删除的映射到Done。
此外,在 Lagom 的 KafkaSubscriberActor 的 impl 中,我们看到 private atLeastOnce 的 impl 基本上解压缩了消息有效负载和偏移量,然后在我们的用户流将消息映射到 Done 后重新压缩然后备份。
上面的这两个花絮似乎暗示,通过使用流管理器和跳过元素,我们最终可能会出现可提交偏移量不再与每个 Kafka 消息生成的Dones 均匀压缩的情况。
示例:如果我们流式传输 1、2、3、4 并将 1、2 和 4 映射到 Done 但在 3 上抛出异常,我们有 3 个Dones 和 4 个可提交的偏移量?
- 这是正确的/预期的吗?这是否意味着我们应该避免在这里使用流监督器?
- 拉链不均匀会导致哪些行为?
- 在通过 Lagom 消息代理 API 使用来自 Kafka 的消息时,推荐的错误处理方法是什么?将故障映射/恢复到
Done是否正确?
使用 Lagom 1.4.10
【问题讨论】:
标签: scala apache-kafka akka lagom