上一篇《Flink Actor源码深度剖析》讲了 Flink 中 Actor 模型的应用和 Akka RPC 框架的源码实现。但很多人看完后仍然有疑问ActorSystem 内部到底是怎么管理 Actor 的一条消息从发送到接收底层经历了哪些步骤Dispatcher 调度器是如何将 Actor 绑定到线程执行的这篇深入 Akka 底层从 ActorSystem 架构、Actor 内部结构、消息传递机制、Dispatcher 调度器四个维度把 Akka 的底层实现讲透。一、Akka ActorSystem底层架构下面这张图是 Akka ActorSystem 底层架构包括 Actor 层级、Actor 内部结构和 Actor 路径与寻址。1.1 ActorSystem是什么ActorSystem 是 Akka 的核心入口负责管理所有 Actor 的生命周期。可以把 ActorSystem 理解为一个Actor 的世界所有 Actor 都在这个世界中创建、运行和销毁。ActorSystem 的核心职责配置管理加载和管理 Akka 配置application.conf调度器管理创建和管理 Dispatcher 调度器和线程池事件流管理 EventStream用于发布订阅事件日志系统管理 LoggingBus 和日志适配器扩展机制管理 Akka Extension如 Cluster、Persistence、RemotingActor 创建通过 actorOf() 创建顶级 Actor创建 ActorSystem 的代码// 创建默认配置的 ActorSystemActorSystemsystemActorSystem.create(flink);// 使用自定义配置ConfigconfigConfigFactory.load(akka.conf);ActorSystemsystemActorSystem.create(flink,config);1.2 Actor层级结构Akka 中的 Actor 天然形成层级结构每个 Actor 都有一个父 Actor最顶层是三个 Guardian Actor守护者Root Guardian/所有 Actor 的根监督 System Guardian 和 User GuardianSystem Guardian/system监督系统级 Actor如日志、事件流、远程传输等User Guardian/user监督用户创建的顶级 Actor通过 system.actorOf() 创建Actor 层级的创建方式// 创建顶级 Actor在 /user 下ActorRefparentsystem.actorOf(Props.create(ParentActor.class),parent);// 在 Actor 内部创建子 Actor在 /user/parent 下ActorRefchildgetContext().actorOf(Props.create(ChildActor.class),child);Actor 路径示例akka://flink/user/parent— 顶级 Actorakka://flink/user/parent/child— 子 Actorakka://flink/system/log1— 系统 Actor1.3 Actor内部结构很多人以为 Actor 就是一个简单的对象实际上 Actor 的内部结构非常精巧由四层组成第一层ActorRefActor 引用ActorRef 是对外的唯一引用用于发送消息隐藏 Actor 的内部实现支持本地和远程透明主要实现LocalActorRef本地、RemoteActorRef远程、RepointableActorRef可重定位第二层ActorCellActor 单元ActorCell 是 Actor 的核心管理单元持有 Actor 实例、Mailbox、Dispatcher、父 Actor 引用、子 Actor 列表负责消息分发、生命周期管理、监督处理ActorRef.tell() 最终调用的是 ActorCell.sendMessage()第三层Mailbox邮箱Mailbox 是消息队列FIFO 顺序实现了 Runnable 接口可以被线程池执行内部包含 MessageQueue实际存储消息的队列默认实现UnboundedMailbox无界、BoundedMailbox有界、PriorityMailbox优先级第四层ActorActor 实例Actor 是实际处理消息的对象包含 receive() 方法通过模式匹配处理消息持有 contextActorContext可以创建子 Actor、监督、获取自身引用用户自定义的 Actor 逻辑都在这里Actor 内部结构的关系ActorRef → ActorCell → Mailbox → Actor (引用) (管理单元) (消息队列) (业务逻辑)1.4 Actor路径与寻址Akka 支持多种 Actor 引用类型实现位置透明Location Transparency引用类型说明路径示例LocalActorRef本地 Actor 引用直接调用 ActorCellakka://sys/user/actorRemoteActorRef远程 Actor 引用通过 Akka Remote 发送akka.tcp://syshost:port/user/actorActorSelection通过路径通配符选择多个 Actor/user/worker/*RepointableActorRef可重定位引用支持远程部署时切换本地↔远程透明切换位置透明是 Akka 的核心设计理念本地 Actor 和远程 Actor 使用相同的编程模型通过 ActorRef 抽象屏蔽底层传输差异。开发者不需要关心 Actor 是在本地还是远程只需要通过 ActorRef 发送消息即可。二、消息传递底层机制下面这张图是 Akka 消息传递底层机制包括本地消息传递、远程消息传递和 Akka Remote Netty 架构。2.1 本地消息传递6步流程以actorRef.tell(msg, sender)为例本地消息传递的完整流程第1步ActorRef.tell()调用方通过 ActorRef 发送消息指定发送者senderLocalActorRef.tell() 内部调用 underlying.sendMessage()消息和发送者被封装为 Envelope信封第2步ActorCell.sendMessage()ActorCell 接收消息创建 Envelopemessage sender调用 dispatcher.dispatch(this, envelope) 将消息交给 DispatcherActorCell 持有 Dispatcher 的引用第3步MessageDispatcher.dispatch()Dispatcher 将 Envelope 放入 Mailbox 的 MessageQueue调用 Mailbox.setAsScheduled() 标记为已调度通过 executorService.execute(mailbox) 将 Mailbox 提交到线程池第4步Mailbox.run() 被线程调度ExecutorService 从线程池分配一个线程执行 Mailbox 的 run() 方法Mailbox 实现了 Runnable 接口第5步Mailbox.processMailbox()从 MessageQueue 取出消息FIFO 顺序调用 actor.aroundReceive() 或 actor.receive() 处理消息循环处理消息直到队列为空或达到 throughput 限制处理完后如果队列还有消息重新提交到线程池第6步Actor.receive() 处理消息用户自定义的 receive() 方法通过 PartialFunction 模式匹配消息类型执行业务逻辑处理完一条消息后继续取下一条消息关键特性同一 Actor 的消息在同一线程中顺序处理FIFO不同 Actor 可并行处理无需锁和同步。这是 Actor 模型线程安全的基础。2.2 远程消息传递6步流程远程消息传递比本地多了序列化、帧编码、Netty 传输、反序列化等步骤第1步RemoteActorRef.tell()远程 ActorRef 发送消息进入 RemoteTransportRemoteActorRef 内部持有远程地址host:port和 Actor 路径第2步序列化消息通过 Serialization.serialize() 将消息和发送者序列化为字节数组Akka 默认使用 Java 序列化也支持 Protobuf、Kryo 等消息必须实现 Serializable 接口否则序列化失败第3步帧编码将序列化后的字节封装为 Akka 远程帧帧结构帧头协议版本、消息类型、长度 消息体通过 LengthFieldBasedFrameDecoder 解决粘包问题第4步Netty TCP 传输通过 Netty Client 将帧写入 TCP ChannelNetty 的 ChannelPipeline 包含编码器、解码器、业务处理器发送到远端的 Netty Server第5步远端接收反序列化远端 Netty Server 接收字节流通过 LengthFieldBasedFrameDecoder 解码为帧反序列化为消息对象Envelope通过远端 ActorRef 将消息放入目标 Actor 的 Mailbox第6步放入远端 Mailbox后续流程与本地消息传递完全相同Dispatcher 调度 → Mailbox.run() → Actor.receive()位置透明远端 Actor 不需要知道消息来自本地还是远程2.3 Akka Remote Netty架构Akka Remote 底层基于 Netty 实现分为四层第一层Akka 应用层Actor / ActorRef / MessageDispatcher业务逻辑层不关心底层传输第二层Akka Remote 层RemoteTransport远程传输抽象序列化将消息序列化为字节帧编码封装为 Akka 远程帧心跳检测定期发送心跳监控连接状态重连机制连接断开后自动重连握手协议建立连接时的握手交换 UID 和协议版本第三层Netty 层ClientBootstrap / ServerBootstrapNetty 启动器ChannelPipelineChannel 处理流水线ByteBufNetty 字节缓冲区零拷贝LengthFieldBasedFrameDecoder基于长度字段的帧解码器解决粘包MessageEncoder / MessageDecoder消息编码器/解码器第四层TCP 传输层Socket / TCP 连接字节流传输可靠有序传输2.4 序列化机制Akka 远程通信需要序列化消息支持多种序列化方式序列化方式性能体积兼容性适用场景Java 序列化低大好默认兼容所有 SerializableProtobuf高小需定义高性能需定义 .protoKryo高小较好高性能无需定义文件Flink 中 Akka 默认使用 Java 序列化因为 Flink 的 RPC 消息RpcInvocation已经实现了 Serializable。对于性能敏感的场景可以配置使用 Kryo 或 Protobuf。三、Dispatcher调度器底层下面这张图是 Akka Dispatcher 调度器与性能调优包括调度模型、四种调度器类型、配置参数和最佳实践。3.1 Dispatcher调度模型Dispatcher 是 Akka 的核心调度组件负责将 Actor 的 Mailbox 调度到线程池执行。调度模型的完整流程消息 → Mailbox消息队列→ Dispatcher调度注册→ ExecutorService执行器→ ThreadPool线程池→ Actor.receive()消息处理Dispatcher 的核心职责消息分发将消息放入 Mailbox 的 MessageQueue调度注册将 Mailbox 注册到 ExecutorService等待线程执行线程分配从线程池分配线程执行 Mailbox.run()吞吐量控制控制单次调度处理的最大消息数throughputDispatcher 的关键代码逻辑publicclassDispatcherextendsMessageDispatcher{privatefinalExecutorServiceexecutorService;privatefinalintthroughput;Overridepublicvoiddispatch(ActorCellcell,Envelopehandle){// 1. 将消息放入 MailboxMailboxmailboxcell.mailbox();mailbox.enqueue(handle);// 2. 注册 Mailbox 到执行器registerForExecution(mailbox);}protectedvoidregisterForExecution(Mailboxmailbox){if(mailbox.setAsScheduled()){// 3. 提交到线程池执行executorService.execute(mailbox);}}}3.2 四种调度器类型Akka 提供四种调度器适用于不同场景1. Dispatcher默认调度器线程模型共享线程池ForkJoinPool特点事件驱动、非阻塞、多个 Actor 共享线程池适用大多数 ActorCPU 密集型和非阻塞操作配置akka.actor.default-dispatcher2. PinnedDispatcher独占调度器线程模型每个 Actor 独占一个线程特点为每个 Actor 分配专属线程隔离性好适用阻塞操作同步 IO、Thread.sleep、需要隔离的关键 Actor注意Actor 数量多时线程数爆炸谨慎使用3. CallingThreadDispatcher调用线程调度器线程模型在调用线程中同步执行不创建新线程特点同步执行消息在发送者线程中处理适用仅用于测试不适合生产环境注意会阻塞调用线程4. Custom Dispatcher自定义调度器线程模型自定义线程池配置特点为特定 Actor 配置独立线程池隔离资源适用需要独立资源隔离的 Actor 组如 IO 密集型 Actor配置在 application.conf 中自定义 dispatcher 配置自定义调度器配置示例my-io-dispatcher { type Dispatcher executor thread-pool-executor thread-pool-executor { core-pool-size-min 10 core-pool-size-factor 3.0 core-pool-size-max 30 } throughput 100 }使用自定义调度器ActorRefioActorsystem.actorOf(Props.create(IOActor.class).withDispatcher(my-io-dispatcher),io-actor);3.3 ForkJoinPool底层原理Akka 默认使用 ForkJoinPool 作为线程池这是因为 ForkJoinPool 非常适合 Actor 模型的工作窃取Work Stealing特性。ForkJoinPool 的核心特性工作窃取Work Stealing空闲线程从其他线程的任务队列尾部窃取任务提高线程利用率双端队列Deque每个工作线程有一个双端队列LIFO 处理自己的任务FIFO 窃取其他线程的任务** Fork/Join 任务**支持任务拆分fork和结果合并join适合分治算法轻量级ForkJoinTask 比 Runnable/Callable 更轻量开销更小Actor 模型与 ForkJoinPool 的契合点每个 Mailbox 是一个 ForkJoinTask提交到 ForkJoinPool空闲线程可以窃取其他 Mailbox 任务避免线程空闲Actor 消息处理是轻量级的适合 ForkJoinTaskForkJoinPool 并行度配置fork-join-executor { parallelism-min 8 # 最小线程数 parallelism-factor 2.0 # 并行度因子线程数CPU核数×因子 parallelism-max 64 # 最大线程数 }实际线程数 max(min, min(max, CPU核数 × factor))3.4 Mailbox调度机制Mailbox 是消息队列同时实现了 Runnable 接口可以被线程池执行。Mailbox 的调度机制Mailbox 的状态机Idle空闲队列为空未被调度Scheduled已调度已提交到线程池等待执行Running运行中正在被线程执行处理消息Mailbox.run() 的执行逻辑publicclassMailboximplementsRunnable{privatefinalMessageQueuemessageQueue;privatefinalActorcell;privatevolatilebooleanscheduledfalse;Overridepublicvoidrun(){try{// 1. 处理消息最多处理 throughput 条intprocessed0;while(processedthroughput!messageQueue.isEmpty()){EnvelopemsgmessageQueue.dequeue();cell.aroundReceive(msg.message(),msg.sender());processed;}}finally{// 2. 处理完后如果队列还有消息重新调度setAsIdle();if(!messageQueue.isEmpty()){dispatcher.registerForExecution(this);}}}publicbooleansetAsScheduled(){if(!scheduled){scheduledtrue;returntrue;}returnfalse;// 已调度避免重复提交}}关键设计setAsScheduled()原子操作避免 Mailbox 被重复提交到线程池throughput 控制单次调度最多处理 throughput 条消息避免一个 Actor 长时间占用线程重新调度处理完后如果队列还有消息重新提交到线程池保证消息不丢失3.5 配置参数详解Akka Dispatcher 的关键配置参数参数默认值说明parallelism-factor2.0并行度因子线程数CPU核数×因子parallelism-min8最小线程数parallelism-max64最大线程数core-pool-size-min8线程池核心线程数最小值ThreadPoolExecutorcore-pool-size-factor3.0核心线程数因子ThreadPoolExecutorcore-pool-size-max64核心线程数最大值ThreadPoolExecutorthroughput5单次调度处理的最大消息数throughput-deadline-time0ms吞吐量截止时间0表示无截止mailbox-capacity1000邮箱容量-1表示无界mailbox-typeUnboundedMailbox邮箱类型fairnessfalse是否公平调度 Mailboxthroughput 参数的影响throughput 大如 100减少调度开销提高吞吐量但增加单个 Actor 的延迟throughput 小如 1降低延迟提高响应性但增加调度开销默认值 5 是吞吐量和延迟的平衡点mailbox-capacity 的影响无界邮箱UnboundedMailbox不会丢消息但可能导致 OOM有界邮箱BoundedMailbox容量满后新消息被丢弃或阻塞防止 OOM生产环境建议使用有界邮箱配合监控告警四、Flink中Akka的使用与优化4.1 Flink中Akka的配置Flink 中 Akka 相关的配置参数flink-conf.yaml# Actor 线程池配置akka.actor.default-dispatcher.fork-join-executor.parallelism-factor:2.0akka.actor.default-dispatcher.fork-join-executor.parallelism-min:8akka.actor.default-dispatcher.fork-join-executor.parallelism-max:64# 远程连接配置akka.remote.netty.tcp.connection-timeout:120s# 心跳检测配置akka.remote.watch-failure-detector.heartbeat-interval:10sakka.remote.watch-failure-detector.acceptable-heartbeat-pause:60sakka.remote.transport-failure-detector.heartbeat-interval:10sakka.remote.transport-failure-detector.acceptable-heartbeat-pause:60s# 监督策略akka.actor.guardian-supervisor-strategy:akka.actor.StoppingSupervisorStrategy# 远程事件日志akka.remote.log-remote-lifecycle-events:off# 序列化akka.serialization-java.enabled:on# Akka 框架超时akka.actor.ask-timeout:100sakka.client.timeout:60s4.2 Flink中Akka的优化建议1. 线程池优化JobManagerRPC 调用量不大保持默认配置即可TaskManager如果 Task 数量多可以适当调大 parallelism-max注意Actor 线程池和 Flink 的 Task 执行线程池是分开的不要混淆CPU 密集型parallelism-factor 设为 1.0避免线程切换IO 密集型使用自定义调度器或 PinnedDispatcher 隔离2. 心跳检测优化网络不稳定的环境适当调大 acceptable-heartbeat-pause如 120s避免误判节点故障对延迟敏感的场景调小 heartbeat-interval如 5s更快感知故障跨机房部署必须调大心跳超时避免网络抖动导致节点被误判为故障3. 邮箱优化生产环境建议使用有界邮箱防止 OOM监控邮箱队列长度发现持续增长及时排查高吞吐场景调大 throughput如 10-20减少调度开销低延迟场景 throughput 设为 1尽快处理每条消息4. 序列化优化默认 Java 序列化性能差、体积大性能敏感场景考虑 Kryo 或 Protobuf确保所有 RPC 消息实现 Serializable大对象考虑传引用或分片传输避免单条消息过大5. 避免阻塞操作不要在 Actor 中执行阻塞操作Thread.sleep、同步 IO、Future.get()阻塞操作会占用线程导致其他 Actor 无法被调度IO 密集型 Actor 使用 PinnedDispatcher 或自定义调度器隔离异步操作使用 pipeTo 将 Future 结果回传给 Actor五、总结Flink Akka 底层原理深度剖析要点回顾第一ActorSystem 底层架构是理解 Akka 的基础。ActorSystem 是 Akka 的核心入口管理所有 Actor 的生命周期。Actor 天然形成层级结构Root Guardian → System Guardian → User Guardian → User Actor → Child Actor。Actor 内部由四层组成ActorRef对外引用支持本地/远程透明→ ActorCell管理单元持有 Mailbox/Dispatcher/Actor实例→ Mailbox消息队列FIFO实现 Runnable→ Actor业务逻辑receive() 处理消息。位置透明是 Akka 的核心设计理念通过 ActorRef 抽象屏蔽本地和远程差异。第二消息传递底层机制是 Akka 的核心。本地消息传递 6 步流程ActorRef.tell() → ActorCell.sendMessage() → Dispatcher.dispatch() → Mailbox.run() 被线程调度 → Mailbox.processMailbox() → Actor.receive()。远程消息传递多了序列化、帧编码、Netty TCP 传输、反序列化等步骤。Akka Remote 基于 Netty 实现分为应用层、Remote 层、Netty 层、TCP 层四层。关键特性同一 Actor 的消息顺序处理不同 Actor 并行处理无需锁和同步。第三Dispatcher 调度器是 Akka 性能的关键。Dispatcher 负责将 Mailbox 调度到线程池执行调度模型消息 → Mailbox → Dispatcher → ExecutorService → ThreadPool → Actor.receive()。四种调度器Dispatcher默认共享 ForkJoinPool、PinnedDispatcher独占线程适合阻塞操作、CallingThreadDispatcher调用线程仅用于测试、Custom Dispatcher自定义线程池资源隔离。ForkJoinPool 的工作窃取特性非常适合 Actor 模型。Mailbox 实现 Runnable通过 setAsScheduled() 避免重复调度throughput 控制单次处理消息数。第四配置参数与性能调优是生产环境的关键。核心参数parallelism-factor/min/max线程池大小、throughput单次处理消息数、mailbox-capacity邮箱容量、heartbeat-interval/acceptable-heartbeat-pause心跳检测。调优建议CPU 密集型 parallelism-factor1.0IO 密集型用 PinnedDispatcher 隔离高吞吐调大 throughput低延迟 throughput1生产环境用有界邮箱防 OOM避免在 Actor 中执行阻塞操作。第五Flink 中 Akka 的使用与优化Flink 使用 Akka 作为 RPC 框架JobManager 和 TaskManager 各有独立的 ActorSystem。优化建议包括线程池配置、心跳检测调优、邮箱优化、序列化优化、避免阻塞操作等。Flink 2.0 已移除 Akka 依赖改用自研 RPC 框架但 Actor 模型的设计思想仍然值得深入学习。Akka 的底层设计体现了分布式系统的经典设计思想位置透明、消息驱动、监督容错、工作窃取。理解这些底层原理不仅能帮助排查 Flink RPC 问题更能体会到分布式系统设计的精妙之处。
阅读完成 · 觉得有帮助?