Akka(35): Http:Server side streaming

2024-04-09 04:48
文章标签 http server 35 akka streaming side

本文主要是介绍Akka(35): Http:Server side streaming,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

   在前面几篇讨论里我们都提到过:Akka-http是一项系统集成工具库。它是以数据交换的形式进行系统集成的。所以,Akka-http的核心功能应该是数据交换的实现了:应该能通过某种公开的数据格式和传输标准比较方便的实现包括异类系统之间通过网上进行的数据交换。覆盖包括:数据编码、发送和数据接收、解析全过程。Akka-http提供了许多网上传输标准数据的概括模型以及数据类型转换方法,可以使编程人员很方便的构建网上往来的Request和Response。但是,现实中的数据交换远远不止针对request和response操作能够满足的。系统之间数据交换经常涉及文件或者数据库表类型的数据上传下载。虽然在Http标准中描述了如何通过MultiPart消息类型进行批量数据的传输,但是这个标准涉及的实现细节包括数据内容描述、数据分段方式、消息数据长度计算等等简直可以立即令人却步。Akka-http是基于Akka-stream开发的:不但它的工作流程可以用Akka-stream来表达,它还支持stream化的数据传输。我们知道:Akka-stream提供了功能强大的FileIO和Data-Streaming,可以用Stream-Source代表文件或数据库数据源。简单来说:Akka-http的消息数据内容HttpEntity可以支持理论上无限长度的data-stream。最可贵的是:这个Source是个Reactive-Stream-Source,具备了back-pressure机制,可以有效应付数据交换参与两方Reactive端点不同的数据传输速率。

  Akka-http的stream类型数据内容是以Source[T,_]类型表示的。首先,Akka-stream通过FileIO对象提供了足够多的file-io操作函数,其中有个fromPath函数可以用某个文件内容数据构建一个Source类型:

/*** Creates a Source from a files contents.* Emitted elements are `chunkSize` sized [[akka.util.ByteString]] elements,* except the final element, which will be up to `chunkSize` in size.** You can configure the default dispatcher for this Source by changing the `akka.stream.blocking-io-dispatcher` or* set it for a given Source by using [[akka.stream.ActorAttributes]].** It materializes a [[Future]] of [[IOResult]] containing the number of bytes read from the source file upon completion,* and a possible exception if IO operation was not completed successfully.** @param f         the file path to read from* @param chunkSize the size of each read operation, defaults to 8192*/def fromPath(f: Path, chunkSize: Int = 8192): Source[ByteString, Future[IOResult]] =fromPath(f, chunkSize, startPosition = 0)

这个函数构建了Source[ByteString,Future[IOResult]],我们需要把ByteString转化成MessageEntity。首先需要在implicit-scope内提供Marshaller[ByteString,MessageEntity]类型的隐式实例:

trait JsonCodec extends Json4sSupport {import org.json4s.DefaultFormatsimport org.json4s.ext.JodaTimeSerializersimplicit val serilizer = jackson.Serializationimplicit val formats = DefaultFormats ++ JodaTimeSerializers.all
}
object JsConverters extends JsonCodecobject ServerStreaming extends App {import JsConverters._
...

我们还需要Json-Streaming支持:

  implicit val jsonStreamingSupport = EntityStreamingSupport.json().withParallelMarshalling(parallelism = 8, unordered = false)

FileIO是blocking操作,我们还可以选用独立的线程供blocking操作使用:

   FileIO.fromPath(file, 256).withAttributes(ActorAttributes.dispatcher("akka.http.blocking-ops-dispatcher"))

现在我们可以从在server上用一个文件构建Source然后再转成Response:

  val route =get {path("files"/Remaining) { name =>complete(loadFile(name))} }def loadFile(path: String) = {//   implicit val ec = httpSys.dispatchers.lookup("akka.http.blocking-ops-dispatcher")val file = Paths.get("/Users/tiger/"+path)FileIO.fromPath(file, 256).withAttributes(ActorAttributes.dispatcher("akka.http.blocking-ops-dispatcher")).map(_.utf8String)}

同样,我们也可以把数据库表内数据转成Akka-Stream-Source,然后再实现到MessageEntity的转换。转换过程包括用Query读取数据库表内数据后转成Reactive-Publisher,然后把publisher转成Akka-Stream-Source,如下:

object SlickDAO {import slick.jdbc.H2Profile.api._val dbConfig: slick.basic.DatabaseConfig[slick.jdbc.H2Profile] = slick.basic.DatabaseConfig.forConfig("slick.h2")val db = dbConfig.dbcase class CountyModel(id: Int, name: String)case class CountyTable(tag: Tag) extends Table[CountyModel](tag,"COUNTY") {def id = column[Int]("ID",O.AutoInc,O.PrimaryKey)def name = column[String]("NAME",O.Length(64))def * = (id,name)<>(CountyModel.tupled,CountyModel.unapply)}val CountyQuery = TableQuery[CountyTable]def loadTable(filter: String) = {//   implicit val ec = httpSys.dispatchers.lookup("akka.http.blocking-ops-dispatcher")val qry = CountyQuery.filter {_.name.toUpperCase like s"%${filter.toUpperCase}%"}val publisher = db.stream(qry.result)Source.fromPublisher(publisher = publisher).withAttributes(ActorAttributes.dispatcher("akka.http.blocking-ops-dispatcher"))}
}

然后进行到MessageEntity的转换:

  val route =get {path("files"/Remaining) { name =>complete(loadFile(name))} ~path("tables"/Segment) { t =>complete(SlickDAO.loadTable(t))}}

下面是本次示范的完整源代码:

import java.nio.file._
import akka.actor._
import akka.stream._
import akka.stream.scaladsl._
import akka.http.scaladsl.Http
import akka.http.scaladsl.server.Directives._
import akka.http.scaladsl.common._
import de.heikoseeberger.akkahttpjson4s.Json4sSupport
import org.json4s.jacksonobject SlickDAO {import slick.jdbc.H2Profile.api._val dbConfig: slick.basic.DatabaseConfig[slick.jdbc.H2Profile] = slick.basic.DatabaseConfig.forConfig("slick.h2")val db = dbConfig.dbcase class CountyModel(id: Int, name: String)case class CountyTable(tag: Tag) extends Table[CountyModel](tag,"COUNTY") {def id = column[Int]("ID",O.AutoInc,O.PrimaryKey)def name = column[String]("NAME",O.Length(64))def * = (id,name)<>(CountyModel.tupled,CountyModel.unapply)}val CountyQuery = TableQuery[CountyTable]def loadTable(filter: String) = {//   implicit val ec = httpSys.dispatchers.lookup("akka.http.blocking-ops-dispatcher")val qry = CountyQuery.filter {_.name.toUpperCase like s"%${filter.toUpperCase}%"}val publisher = db.stream(qry.result)Source.fromPublisher(publisher = publisher).withAttributes(ActorAttributes.dispatcher("akka.http.blocking-ops-dispatcher"))}
}trait JsonCodec extends Json4sSupport {import org.json4s.DefaultFormatsimport org.json4s.ext.JodaTimeSerializersimplicit val serilizer = jackson.Serializationimplicit val formats = DefaultFormats ++ JodaTimeSerializers.all
}
object JsConverters extends JsonCodecobject ServerStreaming extends App {import JsConverters._implicit val httpSys = ActorSystem("httpSystem")implicit val httpMat = ActorMaterializer()implicit val httpEC = httpSys.dispatcherimplicit val jsonStreamingSupport = EntityStreamingSupport.json().withParallelMarshalling(parallelism = 8, unordered = false)val (port, host) = (8011,"localhost")val route =get {path("files"/Remaining) { name =>complete(loadFile(name))} ~path("tables"/Segment) { t =>complete(SlickDAO.loadTable(t))}}def loadFile(path: String) = {//   implicit val ec = httpSys.dispatchers.lookup("akka.http.blocking-ops-dispatcher")val file = Paths.get("/Users/tiger/"+path)FileIO.fromPath(file, 256).withAttributes(ActorAttributes.dispatcher("akka.http.blocking-ops-dispatcher")).map(_.utf8String)}val bindingFuture = Http().bindAndHandle(route,host,port)println(s"Server running at $host $port. Press any key to exit ...")scala.io.StdIn.readLine()bindingFuture.flatMap(_.unbind()).onComplete(_ => httpSys.terminate())}

resource/application.conf

akka {http {blocking-ops-dispatcher {type = Dispatcherexecutor = "thread-pool-executor"thread-pool-executor {// or in Akka 2.4.2+fixed-pool-size = 16}throughput = 100}}
}
slick {h2 {driver = "slick.driver.H2Driver$"db {url = "jdbc:h2:~/slickdemo;mv_store=false"driver = "org.h2.Driver"connectionPool = HikariCPnumThreads = 48maxConnections = 48minConnections = 12keepAliveConnection = true}}
}






这篇关于Akka(35): Http:Server side streaming的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/887203

相关文章

SQL Server修改数据库名及物理数据文件名操作步骤

《SQLServer修改数据库名及物理数据文件名操作步骤》在SQLServer中重命名数据库是一个常见的操作,但需要确保用户具有足够的权限来执行此操作,:本文主要介绍SQLServer修改数据... 目录一、背景介绍二、操作步骤2.1 设置为单用户模式(断开连接)2.2 修改数据库名称2.3 查找逻辑文件名

SQL Server数据库死锁处理超详细攻略

《SQLServer数据库死锁处理超详细攻略》SQLServer作为主流数据库管理系统,在高并发场景下可能面临死锁问题,影响系统性能和稳定性,这篇文章主要给大家介绍了关于SQLServer数据库死... 目录一、引言二、查询 Sqlserver 中造成死锁的 SPID三、用内置函数查询执行信息1. sp_w

Maven 配置中的 <mirror>绕过 HTTP 阻断机制的方法

《Maven配置中的<mirror>绕过HTTP阻断机制的方法》:本文主要介绍Maven配置中的<mirror>绕过HTTP阻断机制的方法,本文给大家分享问题原因及解决方案,感兴趣的朋友一... 目录一、问题场景:升级 Maven 后构建失败二、解决方案:通过 <mirror> 配置覆盖默认行为1. 配置示

Linux中修改Apache HTTP Server(httpd)默认端口的完整指南

《Linux中修改ApacheHTTPServer(httpd)默认端口的完整指南》ApacheHTTPServer(简称httpd)是Linux系统中最常用的Web服务器之一,本文将详细介绍如何... 目录一、修改 httpd 默认端口的步骤1. 查找 httpd 配置文件路径2. 编辑配置文件3. 保存

Windows Server 2025 搭建NPS-Radius服务器的步骤

《WindowsServer2025搭建NPS-Radius服务器的步骤》本文主要介绍了通过微软的NPS角色实现一个Radius服务器,身份验证和证书使用微软ADCS、ADDS,具有一定的参考价... 目录简介示意图什么是 802.1X?核心作用802.1X的组成角色工作流程简述802.1X常见应用802.

C++ HTTP框架推荐(特点及优势)

《C++HTTP框架推荐(特点及优势)》:本文主要介绍C++HTTP框架推荐的相关资料,本文通过实例代码给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录1. Crow2. Drogon3. Pistache4. cpp-httplib5. Beast (Boos

SQL Server身份验证模式步骤和示例代码

《SQLServer身份验证模式步骤和示例代码》SQLServer是一个广泛使用的关系数据库管理系统,通常使用两种身份验证模式:Windows身份验证和SQLServer身份验证,本文将详细介绍身份... 目录身份验证方式的概念更改身份验证方式的步骤方法一:使用SQL Server Management S

Spring AI 实现 STDIO和SSE MCP Server的过程详解

《SpringAI实现STDIO和SSEMCPServer的过程详解》STDIO方式是基于进程间通信,MCPClient和MCPServer运行在同一主机,主要用于本地集成、命令行工具等场景... 目录Spring AI 实现 STDIO和SSE MCP Server1.新建Spring Boot项目2.a

SQL Server中的PIVOT与UNPIVOT用法具体示例详解

《SQLServer中的PIVOT与UNPIVOT用法具体示例详解》这篇文章主要给大家介绍了关于SQLServer中的PIVOT与UNPIVOT用法的具体示例,SQLServer中PIVOT和U... 目录引言一、PIVOT:将行转换为列核心作用语法结构实战示例二、UNPIVOT:将列编程转换为行核心作用语

SpringBoot中HTTP连接池的配置与优化

《SpringBoot中HTTP连接池的配置与优化》这篇文章主要为大家详细介绍了SpringBoot中HTTP连接池的配置与优化的相关知识,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一... 目录一、HTTP连接池的核心价值二、Spring Boot集成方案方案1:Apache HttpCl