- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
导读
本文以 Akka 官方文档 io-tcp.md 为主体,系统讲解如何通过Classic Actor(经典 Actor API)使用 Akka 提供的底层 TCP 传输能力。你将掌握:如何获取 TCP Manager、建立出站连接(Connect)、监听入站连接(Bind)、关闭连接(Close / ConfirmedClose / Abort)、通过WriteCommand族写入数据,以及最重要的 ACK/NACK 写入回压与 Push/Pull 读取回压模型,并辅以仓库内完整可运行的 Echo 服务器源码作为实战范本。
1. 依赖引入
TCP I/O API 属于akka-actor模块(所有 Akka I/O API 均内置于该模块,无需额外引入独立模块)。在项目中添加如下依赖(以 Akka BOM 统一管理版本):
// sbt libraryDependencies += "com.typesafe.akka" %% "akka-actor" % AkkaVersion<!-- Maven --> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-actor_2.13</artifactId> <version>AkkaVersion</version> </dependency>// Gradle implementation 'com.typesafe.akka:akka-actor_2.13:AkkaVersion'注意:Akka 依赖托管于 Akka 官方安全库仓库(secure library repository),需要按 https://account.akka.io/token 中的说明配置带 token 的 URL 才能访问。
2. 总体模型:一切经由 Manager Actor
Akka I/O 的所有 API 都通过manager 对象访问。使用任何 I/O API 的第一步,就是获取对应 manager 的引用。TCP 的 manager 是akka.io.Tcp(扩展ExtensionId[TcpExt]),获取方式如下:
Scala
import akka.io.{ IO, Tcp } import context.system // 隐式提供给 IO(Tcp) val manager = IO(Tcp)Java
final ActorRef tcpManager = Tcp.get(getContext().getSystem()).manager();(完整示例见 IODocSpec.scala 与 EchoManager.java。)
Manager 的角色:它本身是一个 Actor,负责管理底层 I/O 资源(selector、channel),并为具体任务实例化 worker——例如为监听入站连接创建专门的监听器。因此,与 TCP 的交互本质上都是“向 manager 发送命令消息、由 manager 派生的子 actor 回复事件消息”的纯 Actor 消息驱动过程。
提示:所有 TCP 相关的命令/事件消息均定义于 Tcp.scala,Scala 直接使用
Tcp.XXX,Java 通过TcpMessage.xxx()工厂方法构造。
3. 出站连接:Connect → Connected → Register
3.1 发送 Connect 命令
连接远程地址的第一步,是向 TCP manager 发送Connect消息:
Scala
IO(Tcp) ! Connect(remote)Java
tcp.tell(TcpMessage.connect(remote), getSelf());Connect消息除最简单的形式外,还支持两个可选参数(源码见 Tcp.scala):
| 参数 | 类型 | 说明 |
|---|---|---|
localAddress | Option[InetSocketAddress]/InetSocketAddress | 指定本地要绑定的地址;不指定则由内核自动选择 |
options | Iterable[SocketOption] | 套接字选项列表,如SO.TcpNoDelay |
timeout | Option[FiniteDuration]/Duration | 连接超时 |
pullMode | Boolean(默认false) | 是否启用 Pull 读取模式,详见第 8 节 |
SO_NODELAY 默认开启:在 Akka 中,SO_NODELAY(Windows 上为TCP_NODELAY)套接字选项默认即为true,与操作系统默认值无关。这会禁用 Nagle 算法,显著降低大多数应用中的延迟。如需覆盖,可在 options 列表中传入SO.TcpNoDelay(false)。
3.2 连接建立流程
TCP manager 收到Connect后,要么回复CommandFailed,要么派生一个内部 actor 代表这条新连接。这个连接 actor 会向原 Connect 命令的发送者发送Connected消息。
连接在被激活之前无法使用:必须向连接 actor 发送Register消息,告诉它“谁将接收来自套接字的数据”。若迟迟不注册,连接 actor 内部存在一个超时,超时后会自动关闭自身并清理资源。
完整客户端示例(Scala):
class Client(remote: InetSocketAddress, listener: ActorRef) extends Actor { import Tcp._ import context.system IO(Tcp) ! Connect(remote) def receive = { case CommandFailed(_: Connect) => listener ! "connect failed" context.stop(self) case c @ Connected(remote, local) => listener ! c val connection = sender() connection ! Register(self) context.become { case data: ByteString => connection ! Write(data) case CommandFailed(_: Write) => // O/S buffer was full listener ! "write failed" case Received(data) => listener ! data case "close" => connection ! Close case _: ConnectionClosed => listener ! "connection closed" context.stop(self) } } }Java 对应版本使用TcpMessage.register(getSelf())激活连接,并用getContext().become(connected(...))切换到已连接状态(完整代码见 IODocTest.java)。
3.3 Register 的语义
Register消息(Tcp.scala)携带三个字段:
final case class Register( handler: ActorRef, keepOpenOnPeerClosed: Boolean = false, useResumeWriting: Boolean = true)handler:接收所有入站数据、并获知连接关闭事件的 actor;keepOpenOnPeerClosed:为true时,对端关闭其写半端后连接不会自动关闭,需显式发出关闭命令(用于实现半关闭,详见第 5 节);useResumeWriting:为false时启用NACK 模式写入回压(不暂停写入,失败即回复CommandFailed),为true时启用带挂起的 NACK 模式(详见第 7 节)。
另外,连接 actor 会watch(监视)注册的 handler:当 handler 终止时,连接会被关闭并释放全部内部资源。
4. 入站连接:Bind → Bound → Connected
4.1 发起绑定
创建 TCP 服务器并监听入站连接,需要向 TCP manager 发送Bind命令:
Scala
IO(Tcp) ! Bind(self, new InetSocketAddress("localhost", 0))Java
tcp.tell(TcpMessage.bind(getSelf(), new InetSocketAddress("localhost", 0), 100), getSelf());Bind消息的参数(Tcp.scala):
| 参数 | 默认值 | 说明 |
|---|---|---|
handler | — | 接收所有Connected消息的 actor |
localAddress | — | 监听的地址;端口指定为0表示绑定随机端口,实际端口见Bound消息 |
backlog | 100 | 内核为该端口保留的未 accept 连接数上限,超出则拒绝连接 |
options | Nil | 套接字选项 |
pullMode | false | 是否启用 Pull 模式(见第 8.2 节) |
4.2 绑定成功与接受连接
发送Bind的 actor 会收到Bound消息,表示服务器已就绪;Bound携带实际绑定的InetSocketAddress(即解析后的 IP 与正确端口号)。
此后处理连接的方式与出站连接一致:收到Connected后,为每个连接派生一个 handler actor,并将 handler 通过Register注册给连接 actor。写入数据则可由系统中任意 actor向连接 actor(即发送过Connected的 actor)发出。
服务器示例(Scala):
class Server extends Actor { import Tcp._ import context.system IO(Tcp) ! Bind(self, new InetSocketAddress("localhost", 0)) def receive = { case b @ Bound(_) => context.parent ! b case CommandFailed(_: Bind) => context.stop(self) case c @ Connected(_, _) => context.parent ! c val handler = context.actorOf(Props[SimplisticHandler]()) val connection = sender() connection ! Register(handler) } }最简单的 handler:
class SimplisticHandler extends Actor { import Tcp._ def receive = { case Received(data) => sender() ! Write(data) case PeerClosed => context.stop(self) } }该 handler 收到数据即原样写回(Echo),对端关闭则终止自身。
4.3 监听端口的生命周期
出站连接与入站监听的一个关键差异:管理监听端口的内部 actor(即Bound消息的发送者)会watch 监听 actor。当监听 actor 终止时,监听端口随之关闭、相关资源全部释放;但已建立的连接不会因此被终止。
5. 关闭连接
连接可通过向连接 actor 发送三种关闭命令之一来关闭:
5.1 Close(正常关闭)
- 命令:
Close/TcpMessage.close() - 行为:发送
FIN关闭连接,不等待对端确认;待写数据(pending writes)会先被 flush。 - 成功通知:
Closed。
5.2 ConfirmedClose(确认式关闭)
- 命令:
ConfirmedClose/TcpMessage.confirmedClose() - 行为:发送
FIN关闭本端发送方向,但继续接收数据,直到对端也关闭连接为止;待写数据会先被 flush。 - 成功通知:
ConfirmedClosed。
5.3 Abort(立即终止)
- 命令:
Abort/TcpMessage.abort() - 行为:向对端发送
RST立即终止连接;待写数据不会被 flush。 - 成功通知:
Aborted。
5.4 对端关闭与错误关闭事件
PeerClosed:对端关闭连接时发送给监听者(listener)。默认情况下本端随后也会自动关闭连接;若希望在Register中将keepOpenOnPeerClosed设为true,则连接保持打开,直到收到上述关闭命令之一(支持半关闭连接场景)。ErrorClosed:发生错误导致连接被迫关闭时发送给监听者,携带错误原因。
上述关闭通知全部是ConnectionClosed的子类型(源码见 Tcp.scala),不需要细粒度区分关闭事件的监听者可以统一按ConnectionClosed处理。
6. 写入数据:WriteCommand 的三种实现
连接建立后,系统中任意 actor 都可以向连接 actor 发送WriteCommand写入数据。WriteCommand是抽象类,有 3 个具体实现:
6.1 Tcp.Write
最简单的写入命令,包装一个ByteString实例和一个 "ack" 事件:
final case class Write(data: ByteString, ack: Event) extends SimpleWriteCommandByteString是 Akka 的不可变内存数据模型(详见 io.md 的 ByteString 一节),一个或多个数据块的最大总大小为2 GB(2^31 字节)。
6.2 Tcp.WriteFile
直接发送文件中的原始数据:
final case class WriteFile(filePath: String, position: Long, count: Long, ack: Event)它指定磁盘上的一段(连续)字节范围直接经由连接发送,无需先加载进 JVM 内存,因此可"持有"超过 2 GB 的数据,同样支持 ack 事件。对发送大文件场景非常高效。
6.3 Tcp.CompoundWrite
把多个Write和/或WriteFile组合成一个原子写命令一次性写往连接,带来三大好处:
- 最小开销:连接 actor 同一时刻只能处理一个写命令;合并成一个
CompoundWrite可避免用 ACK 协议逐个"喂"给连接 actor; - 原子性保证:
WriteCommand是原子的,合并后其他 actor 无法把写入"插入"到你的写序列中间——多 actor 并发写同一连接时这是很难靠其他手段实现的重要特性; - 子写可独立 ACK:
CompoundWrite的子写本身就是普通Write/WriteFile,各自可以请求 ack;这些 ACK 在对应子写完成时发出。这允许你通过组合一个空的请求 ACK 的写来为一次写入附加多个 ACK,或在任意位置插入中间 ACK 来跟踪CompoundWrite的传输进度。
组合语法:Scala 使用
+:/++:(Tcp.scala),Java 使用prepend系列方法。
7. 写入回压:TCP 连接 actor 的三种背压模式
TCP 连接 actor 的基本模型是没有内部缓冲:同一时刻只能处理一个写命令(即只能缓冲一个尚未完全交给操作系统内核的写)。因此拥塞必须由用户层处理——写入和读取两侧都是如此。
7.1 ACK 模式
每个Write命令携带一个任意对象作为 ack;只要该对象不是NoAck,数据全部成功写入套接字后,ack 对象就会返回给Write的发送者。若在收到此确认前不发起新的写入,就不会因缓冲溢出而失败。
7.2 NACK 模式
前一个写尚未完成时到达的每个写,都会被回复CommandFailed(携带失败的写命令)。仅依赖此机制要求实现的协议能够容忍"跳过"写入(例如每个写本身都是独立有效消息、不要求全部送达)。启用方式:在Register时将useResumeWriting设为false。
7.3 带挂起的 NACK 模式(NACK-based with write suspending)
与 NACK 模式类似,但一旦某个写失败,后续写入将全部失败,直到收到ResumeWriting消息。ResumeWriting会在最后一个已接受写完成后,以WritingResumed消息应答。若驱动连接的 actor 实现了缓冲,并在收到WritingResumed后重发被 NACK 的消息,则每条消息都能恰好一次送达网络套接字。
7.4 读取回压:Push 与 Pull
- Push 读取:连接 actor 一有数据就作为
Received事件发给注册的 reader actor。reader 想向对端施加回压时,发送SuspendReading暂停新数据接收;直到发送ResumeReading后才恢复Received事件。 - Pull 读取:每次发出
Received事件后,连接 actor自动挂起从套接字接收数据,直到 reader 发送ResumeReading。因此新数据通过发送ResumeReading"拉取"而来。详见第 8 节。
重要限制:以上所有流控方案只在"一个写者/读者对一个连接 actor"时成立;多个 actor 同时向同一连接发送写命令时,无法获得一致的结果。
8. 完整实战:ACK 回压 Echo 服务器
以下 Echo 服务器演示 ACK 写入回压 + Push 读取 + 跨 TCP 连接的端到端回压传播。完整源码见 EchoServer.scala(Scala)与 EchoHandler.java、SimpleEchoHandler.java(Java)。
8.1 关键准备:保持半开连接
为保证在关闭连接前能把所有未写数据回写客户端,需要在激活连接时设置keepOpenOnPeerClosed = true:
Scala
case Connected(remote, local) => log.info("received connection from {}", remote) val handler = context.actorOf(Props(handlerClass, sender(), remote)) sender() ! Register(handler, keepOpenOnPeerClosed = true)Java
connection.tell( TcpMessage.register( handler, true, // <-- keepOpenOnPeerClosed flag true), getSelf());服务器管理者
EchoManager使用SupervisorStrategy.stoppingStrategy(连接损坏不可恢复)、preStart中发起Bind、postRestart中直接stop(self)不重启,详见 EchoManager.java 与 EchoServer.scala。
8.2 SimpleEchoHandler:ACK 模式核心逻辑
Scala(SimpleEchoHandler):
class SimpleEchoHandler(connection: ActorRef, remote: InetSocketAddress) extends Actor with ActorLogging { import Tcp._ context.watch(connection) case object Ack extends Event def receive = { case Received(data) => buffer(data) connection ! Write(data, Ack) context.become({ case Received(data) => buffer(data) case Ack => acknowledge() case PeerClosed => closing = true }, discardOld = false) case PeerClosed => context.stop(self) } ... }原则很简单:写完一块数据后,必须等Ack回来才能发送下一块。等待期间切换行为,把新到的数据先缓冲起来。
辅助函数实现缓冲与水位控制:
private def buffer(data: ByteString): Unit = { storage :+= data stored += data.size if (stored > maxStored) { log.warning(s"drop connection to [$remote] (buffer overrun)") context.stop(self) } else if (stored > highWatermark) { log.debug(s"suspending reading") connection ! SuspendReading suspended = true } } private def acknowledge(): Unit = { require(storage.nonEmpty, "storage was empty") val size = storage(0).size stored -= size transferred += size storage = storage.drop(1) if (suspended && stored < lowWatermark) { log.debug("resuming reading") connection ! ResumeReading suspended = false } if (storage.isEmpty) { if (closing) context.stop(self) else context.unbecome() } else connection ! Write(storage(0), Ack) }其中最有趣的是最后一段逻辑:Ack移除缓冲中最旧的数据块;若这是最后一块,则(视对端是否已关闭)关闭连接或回到空闲行为;否则发送下一块缓冲数据并继续等待Ack。同时,缓冲量超过高水位(highWatermark = maxStored * 5 / 10)时发送SuspendReading,低于低水位(lowWatermark)时发送ResumeReading。
8.3 端到端回压传播原理
读取侧回压可以跨连接传播回对端的写者:向连接 actor 发送SuspendReading后,本端不再从套接字读取数据(存在一定延迟,因为连接 actor 处理该命令需要时间,因此缓冲需留出足够余量)。这导致本端 OS 内核缓冲区逐渐填满,进而触发 TCP 窗口机制阻止对端发送、填满对端的写缓冲区,最终对端写者无法再向套接字推送数据——这就是跨 TCP 连接的端到端背压实现机制。
8.4 EchoHandler:NACK 模式 + 写挂起
完整的EchoHandler演示 NACK 模式(源码 EchoServer.scala / EchoHandler.java):
写透传状态(writing):
def writing: Receive = { case Received(data) => connection ! Write(data, Ack(currentOffset)) buffer(data) case Ack(ack) => acknowledge(ack) case CommandFailed(Write(_, Ack(ack))) => connection ! ResumeWriting context.become(buffering(ack)) case PeerClosed => if (storage.isEmpty) context.stop(self) else context.become(closing) }原则是持续写入直到收到CommandFailed,用 ACK 只做重发缓冲的剪枝。收到失败后进入buffering状态处理所有排队数据的重发:
def buffering(nack: Int): Receive = { var toAck = 10 var peerClosed = false { case Received(data) => buffer(data) case WritingResumed => writeFirst() case PeerClosed => peerClosed = true case Ack(ack) if ack < nack => acknowledge(ack) case Ack(ack) => acknowledge(ack) if (storage.nonEmpty) { if (toAck > 0) { writeFirst() // 失败后先用 ACK 模式 toAck -= 1 } else { writeAll() // 然后恢复乐观写透传 context.become(if (peerClosed) closing else writing) } } else if (peerClosed) context.stop(self) else context.become(writing) } }要点:进入该状态时所有已缓冲的写都已发送给连接 actor,因此ResumeWriting排在这些写之后——在收到WritingResumed之前会先收到全部未完成的CommandFailed(此状态下被忽略)。WritingResumed表示内部排队的写已全部完成,后续写入不会失败。EchoHandler据此在失败后先以 ACK 模式写 10 次(toAck = 10),再恢复乐观写透传。
关闭状态(closing):始终发送全部未完成消息、确认所有成功写入;若失败则切换行为等待WritingResumed后重新开始:
def closing: Receive = { case CommandFailed(_: Write) => connection ! ResumeWriting context.become({ case WritingResumed => writeAll() context.unbecome() case ack: Int => acknowledge(ack) }, discardOld = false) case Ack(ack) => acknowledge(ack) if (storage.isEmpty) context.stop(self) }9. Pull 模式读取与入站 Pull 监听
Push 读取模式下,数据一到就发给 actor,因此 Echo 服务器必须为"写入慢于到达"的情况维护入站缓冲。Pull 模式则可以完全消除该缓冲。
9.1 PullEcho:零缓冲回显
Scala(完整代码见 ReadBackPressure.scala):
case object Ack extends Event class PullEcho(connection: ActorRef) extends Actor { override def preStart(): Unit = connection ! ResumeReading def receive = { case Received(data) => connection ! Write(data, Ack) case Ack => connection ! ResumeReading } }Java(见 JavaReadBackPressure.java):
@Override public void preStart() throws Exception { connection.tell(TcpMessage.resumeReading(), getSelf()); } @Override public Receive createReceive() { return receiveBuilder() .match(Tcp.Received.class, message -> { ByteString data = message.data(); connection.tell(TcpMessage.write(data, new Ack()), getSelf()); }) .match(Ack.class, message -> { connection.tell(TcpMessage.resumeReading(), getSelf()); }) .build(); }原理:直到上一次写入被连接 actor 完全 ACK 才恢复读取。每个 Pull 模式连接 actor 初始都处于挂起状态;在preStart中发送ResumeReading通知连接 actor 准备接收第一块数据。因为只在上一块数据完全写完后才恢复读取,无需维护缓冲。
9.2 出站连接的 Pull 模式
将Connect的pullMode参数设为true:
Scala
IO(Tcp) ! Connect(listenAddress, pullMode = true)Java
tcp.tell( TcpMessage.connect( new InetSocketAddress("localhost", 3000), null, options, timeout, true), getSelf());9.3 入站连接的 Pull 模式:Pull 式监听器
将Bind的pullMode参数设为true,则该监听器接受的所有连接都使用 Pull 模式读取:
Scala
override def preStart(): Unit = IO(Tcp) ! Bind(self, new InetSocketAddress("localhost", 0), pullMode = true)Java
tcp.tell( TcpMessage.bind(getSelf(), new InetSocketAddress("localhost", 0), 100, options, true), getSelf());Pull 式监听器的两个额外效应:
- 所有入站连接自动启用 Pull 读取;
- 接受连接本身也变为 Pull 式:处理完一个(或多个)
Connected事件后,监听器必须发送ResumeAccepting才能继续接受新连接。
Pull 模式监听器初始为挂起状态,因此绑定成功后必须先发送ResumeAccepting才开始接受连接:
Scala
def receive = { case Bound(localAddress) => // Accept connections one by one sender() ! ResumeAccepting(batchSize = 1) context.become(listening(sender())) monitor ! localAddress } def listening(listener: ActorRef): Receive = { case Connected(remote, local) => val handler = context.actorOf(Props(classOf[PullEcho], sender())) sender() ! Register(handler, keepOpenOnPeerClosed = true) listener ! ResumeAccepting(batchSize = 1) }Java
.match(Tcp.Bound.class, x -> { listener = getSender(); // Accept connections one by one listener.tell(TcpMessage.resumeAccepting(1), getSelf()); }) .match(Tcp.Connected.class, x -> { ActorRef handler = getContext().actorOf(Props.create(PullEcho.class, getSender())); getSender().tell(TcpMessage.register(handler), getSelf()); // Resume accepting connections listener.tell(TcpMessage.resumeAccepting(1), getSelf()); })如示例所示,处理完一个入站连接后需要再次ResumeAccepting。ResumeAccepting的batchSize参数指定在需要下一次ResumeAccepting之前,可以接受多少个新连接。
10. 小结
Akka 的 TCP I/O API 将底层 socket 操作完全封装进 Actor 模型:
- Manager 单点入口:
IO(Tcp)/Tcp.get(system).manager()获取管理 actor,其余全部通过消息驱动; - 连接生命周期:出站
Connect、入站Bind,成功后以Register激活并注册数据接收者; - 三种关闭语义:
Close(立即 FIN)、ConfirmedClose(半关闭)、Abort(RST),关闭事件统一归入ConnectionClosed; - 写入回压三模式:ACK、NACK、NACK+挂起(
useResumeWriting),CompoundWrite提供原子批量写; - 读取回压两模式:Push(
SuspendReading/ResumeReading)与 Pull(pullMode = true,ResumeReading拉取、监听器ResumeAccepting按批接受连接)。
文中所有示例均取自仓库内可运行测试与文档代码:IODocSpec.scala、IODocTest.java、EchoServer.scala、EchoManager.java、ReadBackPressure.scala 与 JavaReadBackPressure.java,读者可直接对照源码深入研习。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Classic Networking 完全指南:基于 Actor 模型的 TCP/UDP/DNS 底层 I/O 编程
Akka Classic Networking 完全指南:基于 Actor 模型的 TCP/UDP/DNS 底层 I/O 编程 导读 本文围绕 Akka 经典(
后端并发编程异步编程Envoy TCP Proxy 负载卸载点 tcp_proxy_on_data 详解:资源压力下的连接关闭机制
Envoy TCP Proxy 负载卸载点 tcp_proxy_on_data 详解:资源压力下的连接关闭机制 导读 本文围绕 Envoy 在 TCP 代理(T
云原生服务网格网络微服务为什么TCP TIME_WAIT需要等待2MSL?彻底搞懂连接关闭的安全机制
为什么TCP TIME_WAIT需要等待2MSL?彻底搞懂连接关闭的安全机制 TCP连接关闭过程中,TIME_WAIT状态就像一位忠诚的卫士,默默守护着数据传输
文档教程知识库
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考