您可以使用 Kafka 来实现您想要的。 Kafka 主题是持久的(如果您将 Kafka 中的保留设置为 forever 或在主题上启用日志压缩,那么数据将“一直”保留,或者您可以在 Kafka 之外存储偏移量)。
使用 Akka Streams,您将在广播您生成的消息(关于生成主题)之后提交您收到的消息(关于接收主题),从而为您提供“至少一次”传递语义。(对于“exactly-once”,您可以查看 Kafka Transactions)
这是来自 Alpakka Kafka 文档的示例:
Consumer.DrainingControl<Done> control =
Consumer.committableSource(consumerSettings, Subscriptions.topics(topic))
.map(
msg ->
ProducerMessage.single(
new ProducerRecord<>(targetTopic, msg.record().key(), msg.record().value()),
msg.committableOffset() // the passThrough
))
.via(Producer.flexiFlow(producerSettings))
.map(m -> m.passThrough())
.toMat(Committer.sink(committerSettings), Keep.both())
.mapMaterializedValue(Consumer::createDrainingControl)
.run(materializer);
您可以通过几种方式将其与(集群池)Actors 集成。最简单的方法是使用Ask 模式。在这种情况下,流会将消息传递给必须在预定义时间内回复的参与者(可能是self())。当收到回复时,它将在提交原始消息之前在目标流上广播。
这看起来像:
Consumer.DrainingControl<Done> control =
Consumer.committableSource(consumerSettings, Subscriptions.topics(topic))
.mapAsync(1, msg ->
Patterns.ask(actor, msg, Duration.ofSeconds(5))
.thenApply(done ->
ProducerMessage.single(
new ProducerRecord<>(targetTopic, done.key(), done.value()),
msg.committableOffset() // the passThrough
)
)
)
.via(Producer.flexiFlow(producerSettings))
.map(m -> m.passThrough())
.toMat(Committer.sink(committerSettings), Keep.both())
.mapMaterializedValue(Consumer::createDrainingControl)
.run(materializer);
如果您有多个可以同时处理消息的参与者,您还可以增加 mapAsync 调用的并行度。