如何通过Akka Streams Kafka结合WebSocket获取Kafka主题最新及后续消息?
如何用Akka Streams Kafka仅获取Kafka主题的最后一条消息+后续新消息
当然可以实现你的需求!连接WebSocket时只拉取主题的最后一条历史消息,再持续接收新消息,而非所有未读消息,在Akka Streams Kafka里有几种实用方案,我给你详细拆解:
方案一:获取最后一条历史消息+后续新消息
如果需要同时拿到主题的最后一条历史消息,再持续消费新消息,核心思路是先获取每个分区的最新偏移量,手动定位到该偏移量的前一个位置消费最后一条,之后自动跟进新消息:
步骤1:获取分区与最新偏移量
先通过Akka Streams Kafka的工具方法拿到主题的所有分区,以及每个分区的最新结束偏移量:
import org.apache.kafka.common.TopicPartition import akka.kafka.ConsumerSettings import akka.kafka.scaladsl.Consumer val consumerSettings: ConsumerSettings[String, String] = // 你的消费者配置(包含group.id、序列化器等) val topic = "MyTopic" // 获取主题下的所有分区 val partitions = Consumer.partitionsFor(consumerSettings, topic) .map(_.map(p => new TopicPartition(topic, p.partition))) // 获取每个分区的最新结束偏移量 val endOffsets = partitions.flatMap(ps => Consumer.endOffsets(consumerSettings, ps))
步骤2:构造指定偏移量的消费源
基于拿到的偏移量,手动指定每个分区的起始消费位置(最后一条历史消息的位置),再构建消费源:
val kafkaSource = endOffsets.map { offsets => val assignedPartitions = offsets.map { case (tp, endOffset) => // 分区有消息时,定位到最后一条的位置;无消息时从0开始 val startOffset = if (endOffset > 0) endOffset - 1 else 0 tp -> startOffset }.toMap // 创建带初始偏移量的订阅 Subscriptions.assignmentWithOffset(assignedPartitions) }.flatMapConcat(sub => Consumer.plainSource(consumerSettings, sub)) .map(_.value()) // 提取消息内容
整合到WebSocket Flow里
把上面的Kafka源整合到你的WebSocket Flow中:
def source(): Flow[Any, String, NotUsed] = { import akka.stream.scaladsl.Flow // 这里放入上面的consumerSettings、topic、kafkaSource定义 Flow[Any].flatMapConcat(_ => kafkaSource) }
方案二:仅获取后续新消息
如果不需要历史消息,只需要接收连接后的新消息,配置起来更简单:
只需要给消费者设置auto.offset.reset为latest,并且确保你的消费者组是新组(或组内无偏移量记录),这样消费者会直接从最新位置开始消费新消息:
import akka.kafka.ConsumerSettings import akka.kafka.scaladsl.Consumer import org.apache.kafka.clients.consumer.ConsumerConfig def source(): Flow[Any, String, NotUsed] = { import akka.stream.scaladsl.Flow val consumerSettings = ConsumerSettings(system, stringDeserializer, stringDeserializer) .withBootstrapServers("kafka:9092") .withGroupId("my-websocket-group") // 配置从最新位置开始消费 .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest") val kafkaSource = Consumer.plainSource(consumerSettings, Subscriptions.topics("MyTopic")) .map(_.value()) Flow[Any].flatMapConcat(_ => kafkaSource) }
注意事项
- 如果你的消费者组之前已有偏移量记录,
auto.offset.reset不会生效,这时候需要用方案一的手动定位方式,或者先重置消费者组的偏移量。 - 多分区主题下,方案一会自动处理每个分区的最后一条消息,之后持续消费各分区的新消息。
内容的提问来源于stack exchange,提问作者Martin
相关产品推荐
相关产品推荐

