Alpakka JMS Source转Kafka Sink流的消息交付保障咨询
Alpakka JMS Source到Kafka Sink的消息交付保障机制疑问
我搭建了一条Alpakka JMS Source -> Kafka Sink的数据流,查阅Alpakka JMS消费者文档后,仍对该配置下的消息交付保障机制存在疑问。
文档中的示例代码如下:
val result: Future[immutable.Seq[javax.jms.Message]] = jmsSource .take(msgsIn.size) .map { ackEnvelope => ackEnvelope.acknowledge() ackEnvelope.message } .runWith(Sink.seq)
我期望消息仅在Kafka Sink写入成功后才被确认,以此实现**至少一次(at-least-once)**交付保障,但无法仅凭假设确认当前配置是否可行。
考虑到Alpakka未提供跨重启的持久化状态,我认为无法实现类似Flink的**精确一次(exactly-once)**交付,但不确定当前配置是否能保障至少一次交付,还是必须通过Kafka Producer FlexiFlow的map操作来完成消息确认?
内容的提问来源于stack exchange,提问作者Fil Karnicki
相关产品推荐
相关产品推荐

