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

Akka Streams如何实现单生产者多消费者独立线性流?

Akka Streams 内存式单生产者多消费者实现方案

现有组件情况

Alpakka目前没有提供完全匹配Kafka消费组语义的内存式组件,需要自定义实现。

核心实现思路

通过一个中间消息枢纽Actor作为核心,实现单生产者多消费者的消息路由:

  • 生产者通过统一的Sink将消息发送到枢纽Actor
  • 每个消费者通过独立的Source连接到枢纽Actor,维护自身的消费状态(类似Kafka消费组的独立偏移)
  • 枢纽Actor负责消息存储、订阅管理、消息分发及背压处理

具体实现步骤

1. 实现消息枢纽Actor

这个Actor要处理消息存储、消费者订阅、消息拉取和确认逻辑:

import akka.actor.{Actor, ActorRef, Props}

object MessageHubActor {
  case class SendMessage(msg: Any)
  case class Subscribe(consumer: ActorRef)
  case class PullMessages(requested: Int)
  case class ConfirmMessage(msgId: String)
  case class Messages(messages: List[(String, Any)])
}

class MessageHubActor extends Actor {
  import MessageHubActor._

  private var messages: List[(String, Any)] = Nil
  private var consumers: Set[ActorRef] = Set.empty
  private var nextMsgId = 0L

  override def receive: Receive = {
    // 接收生产者发送的消息,分配唯一ID后存入队列并通知所有消费者
    case SendMessage(msg) =>
      val msgId = s"msg-${nextMsgId}"
      nextMsgId += 1
      val msgEntry = (msgId, msg)
      messages = messages :+ msgEntry
      consumers.foreach(_ ! Messages(List(msgEntry)))
    
    // 处理消费者订阅请求,将消费者加入订阅列表并发送当前队列所有消息(可选从头消费)
    case Subscribe(consumer) =>
      consumers += consumer
      if (messages.nonEmpty) consumer ! Messages(messages)
    
    // 处理消费者拉取请求,返回指定数量的消息
    case PullMessages(requested) =>
      val toSend = messages.take(requested)
      sender() ! Messages(toSend)
    
    // 处理消息确认,移除已确认的消息(可根据需求改为按消费者偏移清理)
    case ConfirmMessage(msgId) =>
      messages = messages.filter(_._1 != msgId)
  }
}

2. 封装可复用的Source和Sink

基于枢纽Actor封装成易用的生产者Sink和消费者Source,支持运行时动态创建消费者流:

import akka.stream.scaladsl.{ActorSink, ActorSource, Sink, Source}
import akka.stream.{OverflowStrategy, Materializer}
import akka.actor.ActorSystem

class InMemoryTopic(implicit mat: Materializer, system: ActorSystem) {
  private val hubActor = system.actorOf(Props[MessageHubActor])

  // 生产者Sink:支持背压,将消息发送到枢纽Actor
  def producerSink: Sink[Any, _] = ActorSink.actorRefWithBackpressure(
    ref = hubActor,
    messageAdapter = MessageHubActor.SendMessage,
    onInitMessage = (_: ActorRef) => {}, // 生产者无需订阅,仅发送消息
    ackMessage = MessageHubActor.ConfirmMessage("ack"), // 自定义背压确认信号
    onCompleteMessage = () => {},
    onFailureMessage = (_: Throwable) => {}
  )

  // 消费者Source:物化值为消费者ActorRef,支持动态创建独立消费流
  def consumerSource: Source[(String, Any), ActorRef] = ActorSource.actorRef[(String, Any)](
    completionMatcher = PartialFunction.empty,
    failureMatcher = PartialFunction.empty,
    bufferSize = 100,
    overflowStrategy = OverflowStrategy.dropHead
  ).mapMaterializedValue { consumerActor =>
    hubActor ! MessageHubActor.Subscribe(consumerActor)
    consumerActor
  }
}

3. 使用示例

// 初始化系统和物化器
implicit val system = ActorSystem("InMemoryTopicDemo")
implicit val mat = Materializer(system)

val topic = new InMemoryTopic()

// 生产者流
val producer = Source(List("msg1", "msg2", "msg3")).runWith(topic.producerSink)

// 动态创建消费者1
val consumer1 = topic.consumerSource
  .map { case (id, msg) => s"Consumer1 processed: $msg (id: $id)" }
  .runForeach(println)

// 动态创建消费者2
val consumer2 = topic.consumerSource
  .map { case (id, msg) => s"Consumer2 processed: $msg (id: $id)" }
  .runForeach(println)

方案优势

  • 完全基于内存实现,无外部依赖
  • 每个消费者拥有独立的Source,支持运行时动态创建
  • 消息枢纽Actor天然支持多消费者订阅,类似Kafka多消费组语义
  • 基于Akka Actor的背压机制,保证流处理的稳定性

内容的提问来源于stack exchange,提问作者Alessandro D'Armiento

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 02:30:50