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

Scala+Akka-http中HttpEntity转Protobuf对象的问题排查

Akka HTTP客户端中HttpEntity转Protobuf对象的正确实现方案

问题背景

已基于Scala + Akka HTTP实现了路由中可用的Protobuf编解码器:

trait ProtobufMarshalling[T <: GeneratedMessage, E <: GeneratedMessage] {
  implicit def protobufMarshaller: ToEntityMarshaller[E] = PredefinedToEntityMarshallers.ByteArrayMarshaller.compose[E](r => r.toByteArray)

  implicit def protobufUnmarshaller(implicit companion: GeneratedMessageCompanion[T]): FromEntityUnmarshaller[T] = {
    Unmarshaller.byteArrayUnmarshaller.map[T](bytes => companion.parseFrom(bytes))
  }
}

路由中通过with ProtobufMarshalling正常使用,但作为客户端转换响应的HttpEntity到Protobuf对象时,多次尝试均报错,包括ReadOnlyBufferException、Protobuf解析错误、Unmarshaller找不到等问题。


错误原因分析

  1. ReadOnlyBufferException:StrictEntity.data.asByteBuffer返回只读ByteBuffer,直接调用array()会触发异常,需使用ByteString的原生方法获取字节数组。
  2. Protobuf解析错误(invalid wire type):原entityToBytes的runFold实现可能未正确拼接所有字节,或尝试解析了非预期的响应内容(如错误状态码的响应体)。
  3. Unmarshaller找不到:原特质中的FromEntityUnmarshaller[T]无法直接适配Unmarshal所需的Unmarshaller[ResponseEntity, T],需调整隐式转换的类型定义。

正确实现方案

方案一:复用现有ProtobufMarshalling特质(推荐)

调整特质的Unmarshaller定义,使其适配ResponseEntity,并确保隐式Companion在作用域内:

1. 修改ProtobufMarshalling特质

import akka.http.scaladsl.unmarshalling.Unmarshaller
import akka.http.scaladsl.model.ResponseEntity
import com.google.protobuf.GeneratedMessage
import scalapb.GeneratedMessageCompanion

trait ProtobufMarshalling[T <: GeneratedMessage, E <: GeneratedMessage] {
  // 保持原Marshaller不变
  implicit def protobufMarshaller: ToEntityMarshaller[E] = 
    PredefinedToEntityMarshallers.ByteArrayMarshaller.compose[E](_.toByteArray)

  // 调整Unmarshaller类型为ResponseEntity适配版本
  implicit def protobufUnmarshaller(implicit companion: GeneratedMessageCompanion[T]): Unmarshaller[ResponseEntity, T] = 
    Unmarshaller.byteArrayUnmarshaller.forContentTypes(ContentTypes.`application/octet-stream`)
      .map(companion.parseFrom)
}

2. 在客户端Actor中使用

import akka.http.scaladsl.unmarshalling.Unmarshal
import scala.concurrent.duration._

class ApiFetcher(...) extends Actor with ProtobufMarshalling[ApiResponse, ApiCall] {
  import context.dispatcher
  // 显式提供ApiResponse的Companion隐式实例
  implicit val apiResponseCompanion: GeneratedMessageCompanion[ApiResponse] = ApiResponse
  private val timeout = 10.seconds

  def receive = {
    case HttpResponse(StatusCodes.OK, _, entity, _) =>
      Unmarshal(entity).to[ApiResponse].map { apiResponse =>
        // 处理解析后的ApiResponse
        println(s"Received response: $apiResponse")
      }.recover {
        case ex => 
          println(s"Failed to parse response: ${ex.getMessage}")
      }
    case HttpResponse(status, _, _, _) =>
      println(s"Unexpected response status: $status")
  }
}

方案二:直接实现HttpEntity转字节数组解析

如果不想修改原有特质,可通过toStrict可靠读取响应体,再解析为Protobuf对象:

import akka.http.scaladsl.model.HttpResponse
import scala.concurrent.Future
import scala.concurrent.duration._

private def entityToBytes(entity: ResponseEntity, timeout: FiniteDuration): Future[Array[Byte]] = {
  entity.toStrict(timeout).map(_.data.toArray)
}

// 使用示例
val timeout = 10.seconds
val apiResponse = Await.result(
  responseFuture.flatMap { resp =>
    resp.status match {
      case StatusCodes.OK =>
        entityToBytes(resp.entity, timeout).map(ApiResponse.parseFrom)
      case other =>
        Future.failed(new RuntimeException(s"Request failed with status: $other"))
    }
  },
  timeout
)

关键注意事项

  • 确保响应的Content-Type为application/octet-stream,与服务器端的Protobuf编解码一致。
  • 必须先检查响应状态码,避免解析错误状态下的非Protobuf响应体。
  • 使用Await.result仅适用于测试或同步场景,生产环境建议使用map/flatMap链式调用处理Future。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 00:14:52