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

Alpakka AMQP RPC生产者未收到响应问题排查

问题诊断与修复

你的RPC流程收不到回复,核心问题在于请求与响应的关联机制缺失,以及几个代码细节错误:

关键问题点

  • Correlation ID 未传递:RabbitMQ RPC通过correlationId匹配请求和响应,你的消费者没有把原请求的correlationId附在回复消息上,导致生产者无法识别对应回复。
  • 消息体提取错误:消费者错误地将Envelope作为输入内容处理,而不是从消息中提取实际的消息体字节。
  • 超时设置过短:生产者的ConfirmationTimeout仅200ms,可能消费者还未完成处理就触发了超时。
  • 消费者Source处理不当:你从AmqpIncomingMessage中只提取了envelope和properties,但丢失了消息体,同时回复时没有正确构建带关联ID的属性。

修复后的代码

生产者代码修改

调整超时时间,并确保使用默认的RPC Flow配置(Alpakka会自动处理correlationId和replyTo队列):

import akka.{Done, NotUsed}
import akka.actor.ActorSystem
import akka.stream.ActorMaterializer
import akka.stream.alpakka.amqp.scaladsl._
import akka.stream.alpakka.amqp._
import akka.stream.scaladsl._
import akka.util.ByteString

import scala.concurrent.{ExecutionContext, Future}
import scala.concurrent.duration._

object producer extends App {
  implicit val system = ActorSystem("AmqpRpcExample")
  implicit val materializer = ActorMaterializer()
  implicit val ec: ExecutionContext = system.dispatcher

  val amqpConnectionProvider = AmqpDetailsConnectionProvider("localhost", 5672)
  val queueName = "RPC-QUEUE"
  val queueDeclaration = QueueDeclaration(queueName)

  // 调整超时时间,避免过早超时
  val clientFlow: Flow[WriteMessage, ReadResult, Future[String]] = AmqpRpcFlow.atMostOnceFlow(
    AmqpWriteSettings(amqpConnectionProvider)
      .withRoutingKey(queueName)
      .withDeclaration(queueDeclaration)
      .withBufferSize(10)
      .withConfirmationTimeout(5.seconds), // 延长超时到5秒
    10
  )

  val inputMessages: Source[WriteMessage, NotUsed] = Source(List(
    WriteMessage(ByteString("message 1")),
    WriteMessage(ByteString("message 2")),
    WriteMessage(ByteString("message 3")),
    WriteMessage(ByteString("message 4")),
    WriteMessage(ByteString("message 5")),
  ))

  val responseFutures: Future[Done] = inputMessages
    .via(clientFlow)
    .map(response => {
      val responseStr = response.bytes.utf8String
      println(s"Received response: $responseStr")
      responseStr
    })
    .runWith(Sink.ignore)

  // 等待流完成后关闭系统
  responseFutures.onComplete(_ => system.terminate())
}

消费者代码修改

修复消息体提取,传递correlationId,并正确构建回复消息的属性:

import akka.{Done, NotUsed}
import akka.actor.ActorSystem
import akka.stream.ActorMaterializer
import akka.stream.alpakka.amqp.scaladsl._
import akka.stream.alpakka.amqp._
import akka.stream.scaladsl._
import akka.util.ByteString
import com.rabbitmq.client.AMQP.BasicProperties

import scala.concurrent.{ExecutionContext, Future}
import scala.concurrent.duration._

object consumer extends App {
  implicit val system = ActorSystem("AmqpRpcExample")
  implicit val materializer = ActorMaterializer()
  implicit val ec: ExecutionContext = system.dispatcher

  val amqpConnectionProvider = AmqpDetailsConnectionProvider("localhost", 5672)
  val queueName = "RPC-QUEUE"
  val queueDeclaration = QueueDeclaration(queueName)

  // 正确处理消息:提取消息体,传递correlationId到回复
  val rpcFlow: Flow[AmqpIncomingMessage, WriteMessage, NotUsed] = Flow[AmqpIncomingMessage].map { msg =>
    // 提取实际消息体
    val inputStr = msg.bytes.utf8String
    println(s"Received request: $inputStr")
    
    // 处理请求
    val outputStr = s"Processed: $inputStr"
    
    // 获取原请求的replyTo和correlationId
    val replyTo = msg.properties.getReplyTo
    val correlationId = msg.properties.getCorrelationId
    
    // 构建回复消息的属性,必须携带correlationId
    val replyProps = new BasicProperties.Builder()
      .correlationId(correlationId)
      .build()
    
    // 构建带路由键和属性的回复消息
    WriteMessage(ByteString(outputStr))
      .withRoutingKey(replyTo)
      .withProperties(replyProps)
  }

  val amqpSink: Sink[WriteMessage, Future[Done]] = AmqpSink.replyTo(
    AmqpReplyToSinkSettings(amqpConnectionProvider)
  )

  val amqpSource: Source[AmqpIncomingMessage, NotUsed] = AmqpSource.atMostOnceSource(
    NamedQueueSourceSettings(amqpConnectionProvider, queueName)
      .withDeclaration(queueDeclaration),
    bufferSize = 10
  )

  val rpcServer: Future[Done] = amqpSource
    .via(rpcFlow)
    .runWith(amqpSink)

  println("RPC server started")
  
  // 等待服务终止
  rpcServer.onComplete(_ => system.terminate())
}

核心修复说明

  1. Correlation ID 传递:消费者必须将原请求的correlationId复制到回复消息的属性中,这是Alpakka RpcFlow能够匹配请求和响应的核心标识。
  2. 消息体提取:从AmqpIncomingMessage.bytes中获取实际的请求内容,而不是使用Envelope。
  3. 超时调整:延长生产者的确认超时时间,给消费者足够的处理时间。
  4. 保留完整消息对象:消费者直接使用AmqpIncomingMessage,避免丢失消息体和属性信息。

内容的提问来源于stack exchange,提问作者firas_frikha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 00:42:56