You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 10:02:10