【问题标题】:akka-http send continuous chunked http response (stream)akka-http 发送连续的分块 http 响应(流)
【发布时间】:2016-01-12 10:05:36
【问题描述】:

我有这个带有akka-http 客户端和服务器的粗略测试示例。

Server.scala:

import akka.actor.ActorSystem
import akka.stream.ActorMaterializer
import akka.stream.scaladsl.Sink
import akka.http.scaladsl.Http
import akka.http.scaladsl.model.HttpMethods._
import akka.http.scaladsl.model._
import scala.concurrent.Future

class Server extends Runnable {

    def run() = {

        implicit val system = ActorSystem("server")
        implicit val materializer = ActorMaterializer()

        val serverSource = Http().bind(interface = "localhost", port = 8200)

        val requestHandler: HttpRequest => HttpResponse = {
            case HttpRequest(GET, Uri.Path("/stream"), _, _, _) =>
                HttpResponse(entity = HttpEntity(MediaTypes.`text/plain`, "test"))
        }

        val bindingFuture: Future[Http.ServerBinding] = serverSource.to(Sink.foreach { connection =>
            connection handleWithSyncHandler requestHandler
        }).run()

    }

}

Client.scala:

import akka.actor.ActorSystem
import akka.http.scaladsl.Http
import akka.http.scaladsl.model.{Uri, HttpRequest}
import akka.stream.ActorMaterializer

object Client extends App {

    implicit val system = ActorSystem("client")
    import system.dispatcher

    new Thread(new Server).start()

    implicit val materializer = ActorMaterializer()
    val source = Uri("http://localhost:8200/stream")
    val finished = Http().singleRequest(HttpRequest(uri = source)).flatMap { response =>
        response.entity.dataBytes.runForeach { chunk =>
            println(chunk.utf8String)
        }
    }

}

目前Server 只回复一个“测试”。

如何更改 Server 中的 HttpResponse 以每 1 秒在无限循环中将“测试”作为分块(流)发送?

【问题讨论】:

    标签: scala akka http-streaming akka-stream akka-http


    【解决方案1】:

    找到答案了。

    Server.scala:

    import akka.actor.ActorSystem
    import akka.stream.ActorMaterializer
    import akka.stream.scaladsl.{Source, Sink}
    import akka.http.scaladsl.Http
    import akka.http.scaladsl.model.HttpMethods._
    import akka.http.scaladsl.model._
    import scala.concurrent.Future
    import scala.concurrent.duration._
    
    class Server extends Runnable {
    
        def run() = {
    
            implicit val system = ActorSystem("server")
            implicit val materializer = ActorMaterializer()
    
            val serverSource = Http().bind(interface = "localhost", port = 8200)
    
            val requestHandler: HttpRequest => HttpResponse = {
                case HttpRequest(GET, Uri.Path("/stream"), _, _, _) =>
                    HttpResponse(entity = HttpEntity.Chunked(ContentTypes.`text/plain`, Source(0 seconds, 1 seconds, "test")))
            }
    
            val bindingFuture: Future[Http.ServerBinding] = serverSource.to(Sink.foreach { connection =>
                connection handleWithSyncHandler requestHandler
            }).run()
    
        }
    
    }
    

    【讨论】:

      猜你喜欢
      • 2016-01-21
      • 2018-06-08
      • 2023-04-01
      • 1970-01-01
      • 2018-07-08
      • 2018-07-08
      • 2017-07-17
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多