度器的底層實(shí)現(xiàn))
上一篇《Flink Actor源碼深度剖析》講了 Flink 中 Actor 模型的應(yīng)用和 Akka RPC 框架的源碼實(shí)現(xiàn)。但很多人看完后仍然有疑問(wèn)ActorSystem 內(nèi)部到底是怎么管理 Actor 的一條消息從發(fā)送到接收底層經(jīng)歷了哪些步驟Dispatcher 調(diào)度器是如何將 Actor 綁定到線程執(zhí)行的這篇深入 Akka 底層從 ActorSystem 架構(gòu)、Actor 內(nèi)部結(jié)構(gòu)、消息傳遞機(jī)制、Dispatcher 調(diào)度器四個(gè)維度把 Akka 的底層實(shí)現(xiàn)講透。一、Akka ActorSystem底層架構(gòu)下面這張圖是 Akka ActorSystem 底層架構(gòu)包括 Actor 層級(jí)、Actor 內(nèi)部結(jié)構(gòu)和 Actor 路徑與尋址。1.1 ActorSystem是什么ActorSystem 是 Akka 的核心入口負(fù)責(zé)管理所有 Actor 的生命周期??梢园?ActorSystem 理解為一個(gè)Actor 的世界所有 Actor 都在這個(gè)世界中創(chuàng)建、運(yùn)行和銷(xiāo)毀。ActorSystem 的核心職責(zé)配置管理加載和管理 Akka 配置application.conf調(diào)度器管理創(chuàng)建和管理 Dispatcher 調(diào)度器和線程池事件流管理 EventStream用于發(fā)布訂閱事件日志系統(tǒng)管理 LoggingBus 和日志適配器擴(kuò)展機(jī)制管理 Akka Extension如 Cluster、Persistence、RemotingActor 創(chuàng)建通過(guò) actorOf() 創(chuàng)建頂級(jí) Actor創(chuàng)建 ActorSystem 的代碼// 創(chuàng)建默認(rèn)配置的 ActorSystemActorSystemsystemActorSystem.create(flink);// 使用自定義配置ConfigconfigConfigFactory.load(akka.conf);ActorSystemsystemActorSystem.create(flink,config);1.2 Actor層級(jí)結(jié)構(gòu)Akka 中的 Actor 天然形成層級(jí)結(jié)構(gòu)每個(gè) Actor 都有一個(gè)父 Actor最頂層是三個(gè) Guardian Actor守護(hù)者Root Guardian/所有 Actor 的根監(jiān)督 System Guardian 和 User GuardianSystem Guardian/system監(jiān)督系統(tǒng)級(jí) Actor如日志、事件流、遠(yuǎn)程傳輸?shù)萓ser Guardian/user監(jiān)督用戶(hù)創(chuàng)建的頂級(jí) Actor通過(guò) system.actorOf() 創(chuàng)建Actor 層級(jí)的創(chuàng)建方式// 創(chuàng)建頂級(jí) Actor在 /user 下ActorRefparentsystem.actorOf(Props.create(ParentActor.class),parent);// 在 Actor 內(nèi)部創(chuàng)建子 Actor在 /user/parent 下ActorRefchildgetContext().actorOf(Props.create(ChildActor.class),child);Actor 路徑示例akka://flink/user/parent— 頂級(jí) Actorakka://flink/user/parent/child— 子 Actorakka://flink/system/log1— 系統(tǒng) Actor1.3 Actor內(nèi)部結(jié)構(gòu)很多人以為 Actor 就是一個(gè)簡(jiǎn)單的對(duì)象實(shí)際上 Actor 的內(nèi)部結(jié)構(gòu)非常精巧由四層組成第一層ActorRefActor 引用ActorRef 是對(duì)外的唯一引用用于發(fā)送消息隱藏 Actor 的內(nèi)部實(shí)現(xiàn)支持本地和遠(yuǎn)程透明主要實(shí)現(xiàn)LocalActorRef本地、RemoteActorRef遠(yuǎn)程、RepointableActorRef可重定位第二層ActorCellActor 單元ActorCell 是 Actor 的核心管理單元持有 Actor 實(shí)例、Mailbox、Dispatcher、父 Actor 引用、子 Actor 列表負(fù)責(zé)消息分發(fā)、生命周期管理、監(jiān)督處理ActorRef.tell() 最終調(diào)用的是 ActorCell.sendMessage()第三層Mailbox郵箱Mailbox 是消息隊(duì)列FIFO 順序?qū)崿F(xiàn)了 Runnable 接口可以被線程池執(zhí)行內(nèi)部包含 MessageQueue實(shí)際存儲(chǔ)消息的隊(duì)列默認(rèn)實(shí)現(xiàn)UnboundedMailbox無(wú)界、BoundedMailbox有界、PriorityMailbox優(yōu)先級(jí)第四層ActorActor 實(shí)例Actor 是實(shí)際處理消息的對(duì)象包含 receive() 方法通過(guò)模式匹配處理消息持有 contextActorContext可以創(chuàng)建子 Actor、監(jiān)督、獲取自身引用用戶(hù)自定義的 Actor 邏輯都在這里Actor 內(nèi)部結(jié)構(gòu)的關(guān)系A(chǔ)ctorRef → ActorCell → Mailbox → Actor (引用) (管理單元) (消息隊(duì)列) (業(yè)務(wù)邏輯)1.4 Actor路徑與尋址Akka 支持多種 Actor 引用類(lèi)型實(shí)現(xiàn)位置透明Location Transparency引用類(lèi)型說(shuō)明路徑示例LocalActorRef本地 Actor 引用直接調(diào)用 ActorCellakka://sys/user/actorRemoteActorRef遠(yuǎn)程 Actor 引用通過(guò) Akka Remote 發(fā)送akka.tcp://syshost:port/user/actorActorSelection通過(guò)路徑通配符選擇多個(gè) Actor/user/worker/*RepointableActorRef可重定位引用支持遠(yuǎn)程部署時(shí)切換本地?遠(yuǎn)程透明切換位置透明是 Akka 的核心設(shè)計(jì)理念本地 Actor 和遠(yuǎn)程 Actor 使用相同的編程模型通過(guò) ActorRef 抽象屏蔽底層傳輸差異。開(kāi)發(fā)者不需要關(guān)心 Actor 是在本地還是遠(yuǎn)程只需要通過(guò) ActorRef 發(fā)送消息即可。二、消息傳遞底層機(jī)制下面這張圖是 Akka 消息傳遞底層機(jī)制包括本地消息傳遞、遠(yuǎn)程消息傳遞和 Akka Remote Netty 架構(gòu)。2.1 本地消息傳遞6步流程以actorRef.tell(msg, sender)為例本地消息傳遞的完整流程第1步ActorRef.tell()調(diào)用方通過(guò) ActorRef 發(fā)送消息指定發(fā)送者senderLocalActorRef.tell() 內(nèi)部調(diào)用 underlying.sendMessage()消息和發(fā)送者被封裝為 Envelope信封第2步ActorCell.sendMessage()ActorCell 接收消息創(chuàng)建 Envelopemessage sender調(diào)用 dispatcher.dispatch(this, envelope) 將消息交給 DispatcherActorCell 持有 Dispatcher 的引用第3步MessageDispatcher.dispatch()Dispatcher 將 Envelope 放入 Mailbox 的 MessageQueue調(diào)用 Mailbox.setAsScheduled() 標(biāo)記為已調(diào)度通過(guò) executorService.execute(mailbox) 將 Mailbox 提交到線程池第4步Mailbox.run() 被線程調(diào)度ExecutorService 從線程池分配一個(gè)線程執(zhí)行 Mailbox 的 run() 方法Mailbox 實(shí)現(xiàn)了 Runnable 接口第5步Mailbox.processMailbox()從 MessageQueue 取出消息FIFO 順序調(diào)用 actor.aroundReceive() 或 actor.receive() 處理消息循環(huán)處理消息直到隊(duì)列為空或達(dá)到 throughput 限制處理完后如果隊(duì)列還有消息重新提交到線程池第6步Actor.receive() 處理消息用戶(hù)自定義的 receive() 方法通過(guò) PartialFunction 模式匹配消息類(lèi)型執(zhí)行業(yè)務(wù)邏輯處理完一條消息后繼續(xù)取下一條消息關(guān)鍵特性同一 Actor 的消息在同一線程中順序處理FIFO不同 Actor 可并行處理無(wú)需鎖和同步。這是 Actor 模型線程安全的基礎(chǔ)。2.2 遠(yuǎn)程消息傳遞6步流程遠(yuǎn)程消息傳遞比本地多了序列化、幀編碼、Netty 傳輸、反序列化等步驟第1步RemoteActorRef.tell()遠(yuǎn)程 ActorRef 發(fā)送消息進(jìn)入 RemoteTransportRemoteActorRef 內(nèi)部持有遠(yuǎn)程地址host:port和 Actor 路徑第2步序列化消息通過(guò) Serialization.serialize() 將消息和發(fā)送者序列化為字節(jié)數(shù)組Akka 默認(rèn)使用 Java 序列化也支持 Protobuf、Kryo 等消息必須實(shí)現(xiàn) Serializable 接口否則序列化失敗第3步幀編碼將序列化后的字節(jié)封裝為 Akka 遠(yuǎn)程幀幀結(jié)構(gòu)幀頭協(xié)議版本、消息類(lèi)型、長(zhǎng)度 消息體通過(guò) LengthFieldBasedFrameDecoder 解決粘包問(wèn)題第4步Netty TCP 傳輸通過(guò) Netty Client 將幀寫(xiě)入 TCP ChannelNetty 的 ChannelPipeline 包含編碼器、解碼器、業(yè)務(wù)處理器發(fā)送到遠(yuǎn)端的 Netty Server第5步遠(yuǎn)端接收反序列化遠(yuǎn)端 Netty Server 接收字節(jié)流通過(guò) LengthFieldBasedFrameDecoder 解碼為幀反序列化為消息對(duì)象Envelope通過(guò)遠(yuǎn)端 ActorRef 將消息放入目標(biāo) Actor 的 Mailbox第6步放入遠(yuǎn)端 Mailbox后續(xù)流程與本地消息傳遞完全相同Dispatcher 調(diào)度 → Mailbox.run() → Actor.receive()位置透明遠(yuǎn)端 Actor 不需要知道消息來(lái)自本地還是遠(yuǎn)程2.3 Akka Remote Netty架構(gòu)Akka Remote 底層基于 Netty 實(shí)現(xiàn)分為四層第一層Akka 應(yīng)用層Actor / ActorRef / MessageDispatcher業(yè)務(wù)邏輯層不關(guān)心底層傳輸?shù)诙覣kka Remote 層RemoteTransport遠(yuǎn)程傳輸抽象序列化將消息序列化為字節(jié)幀編碼封裝為 Akka 遠(yuǎn)程幀心跳檢測(cè)定期發(fā)送心跳監(jiān)控連接狀態(tài)重連機(jī)制連接斷開(kāi)后自動(dòng)重連握手協(xié)議建立連接時(shí)的握手交換 UID 和協(xié)議版本第三層Netty 層ClientBootstrap / ServerBootstrapNetty 啟動(dòng)器ChannelPipelineChannel 處理流水線ByteBufNetty 字節(jié)緩沖區(qū)零拷貝LengthFieldBasedFrameDecoder基于長(zhǎng)度字段的幀解碼器解決粘包MessageEncoder / MessageDecoder消息編碼器/解碼器第四層TCP 傳輸層Socket / TCP 連接字節(jié)流傳輸可靠有序傳輸2.4 序列化機(jī)制Akka 遠(yuǎn)程通信需要序列化消息支持多種序列化方式序列化方式性能體積兼容性適用場(chǎng)景Java 序列化低大好默認(rèn)兼容所有 SerializableProtobuf高小需定義高性能需定義 .protoKryo高小較好高性能無(wú)需定義文件Flink 中 Akka 默認(rèn)使用 Java 序列化因?yàn)?Flink 的 RPC 消息RpcInvocation已經(jīng)實(shí)現(xiàn)了 Serializable。對(duì)于性能敏感的場(chǎng)景可以配置使用 Kryo 或 Protobuf。三、Dispatcher調(diào)度器底層下面這張圖是 Akka Dispatcher 調(diào)度器與性能調(diào)優(yōu)包括調(diào)度模型、四種調(diào)度器類(lèi)型、配置參數(shù)和最佳實(shí)踐。3.1 Dispatcher調(diào)度模型Dispatcher 是 Akka 的核心調(diào)度組件負(fù)責(zé)將 Actor 的 Mailbox 調(diào)度到線程池執(zhí)行。調(diào)度模型的完整流程消息 → Mailbox消息隊(duì)列→ Dispatcher調(diào)度注冊(cè)→ ExecutorService執(zhí)行器→ ThreadPool線程池→ Actor.receive()消息處理Dispatcher 的核心職責(zé)消息分發(fā)將消息放入 Mailbox 的 MessageQueue調(diào)度注冊(cè)將 Mailbox 注冊(cè)到 ExecutorService等待線程執(zhí)行線程分配從線程池分配線程執(zhí)行 Mailbox.run()吞吐量控制控制單次調(diào)度處理的最大消息數(shù)throughputDispatcher 的關(guān)鍵代碼邏輯publicclassDispatcherextendsMessageDispatcher{privatefinalExecutorServiceexecutorService;privatefinalintthroughput;Overridepublicvoiddispatch(ActorCellcell,Envelopehandle){// 1. 將消息放入 MailboxMailboxmailboxcell.mailbox();mailbox.enqueue(handle);// 2. 注冊(cè) Mailbox 到執(zhí)行器registerForExecution(mailbox);}protectedvoidregisterForExecution(Mailboxmailbox){if(mailbox.setAsScheduled()){// 3. 提交到線程池執(zhí)行executorService.execute(mailbox);}}}3.2 四種調(diào)度器類(lèi)型Akka 提供四種調(diào)度器適用于不同場(chǎng)景1. Dispatcher默認(rèn)調(diào)度器線程模型共享線程池ForkJoinPool特點(diǎn)事件驅(qū)動(dòng)、非阻塞、多個(gè) Actor 共享線程池適用大多數(shù) ActorCPU 密集型和非阻塞操作配置akka.actor.default-dispatcher2. PinnedDispatcher獨(dú)占調(diào)度器線程模型每個(gè) Actor 獨(dú)占一個(gè)線程特點(diǎn)為每個(gè) Actor 分配專(zhuān)屬線程隔離性好適用阻塞操作同步 IO、Thread.sleep、需要隔離的關(guān)鍵 Actor注意Actor 數(shù)量多時(shí)線程數(shù)爆炸謹(jǐn)慎使用3. CallingThreadDispatcher調(diào)用線程調(diào)度器線程模型在調(diào)用線程中同步執(zhí)行不創(chuàng)建新線程特點(diǎn)同步執(zhí)行消息在發(fā)送者線程中處理適用僅用于測(cè)試不適合生產(chǎn)環(huán)境注意會(huì)阻塞調(diào)用線程4. Custom Dispatcher自定義調(diào)度器線程模型自定義線程池配置特點(diǎn)為特定 Actor 配置獨(dú)立線程池隔離資源適用需要獨(dú)立資源隔離的 Actor 組如 IO 密集型 Actor配置在 application.conf 中自定義 dispatcher 配置自定義調(diào)度器配置示例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 }使用自定義調(diào)度器ActorRefioActorsystem.actorOf(Props.create(IOActor.class).withDispatcher(my-io-dispatcher),io-actor);3.3 ForkJoinPool底層原理Akka 默認(rèn)使用 ForkJoinPool 作為線程池這是因?yàn)?ForkJoinPool 非常適合 Actor 模型的工作竊取Work Stealing特性。ForkJoinPool 的核心特性工作竊取Work Stealing空閑線程從其他線程的任務(wù)隊(duì)列尾部竊取任務(wù)提高線程利用率雙端隊(duì)列Deque每個(gè)工作線程有一個(gè)雙端隊(duì)列LIFO 處理自己的任務(wù)FIFO 竊取其他線程的任務(wù)** Fork/Join 任務(wù)**支持任務(wù)拆分fork和結(jié)果合并join適合分治算法輕量級(jí)ForkJoinTask 比 Runnable/Callable 更輕量開(kāi)銷(xiāo)更小Actor 模型與 ForkJoinPool 的契合點(diǎn)每個(gè) Mailbox 是一個(gè) ForkJoinTask提交到 ForkJoinPool空閑線程可以竊取其他 Mailbox 任務(wù)避免線程空閑Actor 消息處理是輕量級(jí)的適合 ForkJoinTaskForkJoinPool 并行度配置fork-join-executor { parallelism-min 8 # 最小線程數(shù) parallelism-factor 2.0 # 并行度因子線程數(shù)CPU核數(shù)×因子 parallelism-max 64 # 最大線程數(shù) }實(shí)際線程數(shù) max(min, min(max, CPU核數(shù) × factor))3.4 Mailbox調(diào)度機(jī)制Mailbox 是消息隊(duì)列同時(shí)實(shí)現(xiàn)了 Runnable 接口可以被線程池執(zhí)行。Mailbox 的調(diào)度機(jī)制Mailbox 的狀態(tài)機(jī)Idle空閑隊(duì)列為空未被調(diào)度Scheduled已調(diào)度已提交到線程池等待執(zhí)行Running運(yùn)行中正在被線程執(zhí)行處理消息Mailbox.run() 的執(zhí)行邏輯publicclassMailboximplementsRunnable{privatefinalMessageQueuemessageQueue;privatefinalActorcell;privatevolatilebooleanscheduledfalse;Overridepublicvoidrun(){try{// 1. 處理消息最多處理 throughput 條intprocessed0;while(processedthroughput!messageQueue.isEmpty()){EnvelopemsgmessageQueue.dequeue();cell.aroundReceive(msg.message(),msg.sender());processed;}}finally{// 2. 處理完后如果隊(duì)列還有消息重新調(diào)度setAsIdle();if(!messageQueue.isEmpty()){dispatcher.registerForExecution(this);}}}publicbooleansetAsScheduled(){if(!scheduled){scheduledtrue;returntrue;}returnfalse;// 已調(diào)度避免重復(fù)提交}}關(guān)鍵設(shè)計(jì)setAsScheduled()原子操作避免 Mailbox 被重復(fù)提交到線程池throughput 控制單次調(diào)度最多處理 throughput 條消息避免一個(gè) Actor 長(zhǎng)時(shí)間占用線程重新調(diào)度處理完后如果隊(duì)列還有消息重新提交到線程池保證消息不丟失3.5 配置參數(shù)詳解Akka Dispatcher 的關(guān)鍵配置參數(shù)參數(shù)默認(rèn)值說(shuō)明parallelism-factor2.0并行度因子線程數(shù)CPU核數(shù)×因子parallelism-min8最小線程數(shù)parallelism-max64最大線程數(shù)core-pool-size-min8線程池核心線程數(shù)最小值ThreadPoolExecutorcore-pool-size-factor3.0核心線程數(shù)因子ThreadPoolExecutorcore-pool-size-max64核心線程數(shù)最大值ThreadPoolExecutorthroughput5單次調(diào)度處理的最大消息數(shù)throughput-deadline-time0ms吞吐量截止時(shí)間0表示無(wú)截止mailbox-capacity1000郵箱容量-1表示無(wú)界mailbox-typeUnboundedMailbox郵箱類(lèi)型fairnessfalse是否公平調(diào)度 Mailboxthroughput 參數(shù)的影響throughput 大如 100減少調(diào)度開(kāi)銷(xiāo)提高吞吐量但增加單個(gè) Actor 的延遲throughput 小如 1降低延遲提高響應(yīng)性但增加調(diào)度開(kāi)銷(xiāo)默認(rèn)值 5 是吞吐量和延遲的平衡點(diǎn)mailbox-capacity 的影響無(wú)界郵箱UnboundedMailbox不會(huì)丟消息但可能導(dǎo)致 OOM有界郵箱BoundedMailbox容量滿(mǎn)后新消息被丟棄或阻塞防止 OOM生產(chǎn)環(huán)境建議使用有界郵箱配合監(jiān)控告警四、Flink中Akka的使用與優(yōu)化4.1 Flink中Akka的配置Flink 中 Akka 相關(guān)的配置參數(shù)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# 遠(yuǎn)程連接配置akka.remote.netty.tcp.connection-timeout:120s# 心跳檢測(cè)配置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# 監(jiān)督策略akka.actor.guardian-supervisor-strategy:akka.actor.StoppingSupervisorStrategy# 遠(yuǎn)程事件日志akka.remote.log-remote-lifecycle-events:off# 序列化akka.serialization-java.enabled:on# Akka 框架超時(shí)akka.actor.ask-timeout:100sakka.client.timeout:60s4.2 Flink中Akka的優(yōu)化建議1. 線程池優(yōu)化JobManagerRPC 調(diào)用量不大保持默認(rèn)配置即可TaskManager如果 Task 數(shù)量多可以適當(dāng)調(diào)大 parallelism-max注意Actor 線程池和 Flink 的 Task 執(zhí)行線程池是分開(kāi)的不要混淆CPU 密集型parallelism-factor 設(shè)為 1.0避免線程切換IO 密集型使用自定義調(diào)度器或 PinnedDispatcher 隔離2. 心跳檢測(cè)優(yōu)化網(wǎng)絡(luò)不穩(wěn)定的環(huán)境適當(dāng)調(diào)大 acceptable-heartbeat-pause如 120s避免誤判節(jié)點(diǎn)故障對(duì)延遲敏感的場(chǎng)景調(diào)小 heartbeat-interval如 5s更快感知故障跨機(jī)房部署必須調(diào)大心跳超時(shí)避免網(wǎng)絡(luò)抖動(dòng)導(dǎo)致節(jié)點(diǎn)被誤判為故障3. 郵箱優(yōu)化生產(chǎn)環(huán)境建議使用有界郵箱防止 OOM監(jiān)控郵箱隊(duì)列長(zhǎng)度發(fā)現(xiàn)持續(xù)增長(zhǎng)及時(shí)排查高吞吐場(chǎng)景調(diào)大 throughput如 10-20減少調(diào)度開(kāi)銷(xiāo)低延遲場(chǎng)景 throughput 設(shè)為 1盡快處理每條消息4. 序列化優(yōu)化默認(rèn) Java 序列化性能差、體積大性能敏感場(chǎng)景考慮 Kryo 或 Protobuf確保所有 RPC 消息實(shí)現(xiàn) Serializable大對(duì)象考慮傳引用或分片傳輸避免單條消息過(guò)大5. 避免阻塞操作不要在 Actor 中執(zhí)行阻塞操作Thread.sleep、同步 IO、Future.get()阻塞操作會(huì)占用線程導(dǎo)致其他 Actor 無(wú)法被調(diào)度IO 密集型 Actor 使用 PinnedDispatcher 或自定義調(diào)度器隔離異步操作使用 pipeTo 將 Future 結(jié)果回傳給 Actor五、總結(jié)Flink Akka 底層原理深度剖析要點(diǎn)回顧第一ActorSystem 底層架構(gòu)是理解 Akka 的基礎(chǔ)。ActorSystem 是 Akka 的核心入口管理所有 Actor 的生命周期。Actor 天然形成層級(jí)結(jié)構(gòu)Root Guardian → System Guardian → User Guardian → User Actor → Child Actor。Actor 內(nèi)部由四層組成ActorRef對(duì)外引用支持本地/遠(yuǎn)程透明→ ActorCell管理單元持有 Mailbox/Dispatcher/Actor實(shí)例→ Mailbox消息隊(duì)列FIFO實(shí)現(xiàn) Runnable→ Actor業(yè)務(wù)邏輯receive() 處理消息。位置透明是 Akka 的核心設(shè)計(jì)理念通過(guò) ActorRef 抽象屏蔽本地和遠(yuǎn)程差異。第二消息傳遞底層機(jī)制是 Akka 的核心。本地消息傳遞 6 步流程ActorRef.tell() → ActorCell.sendMessage() → Dispatcher.dispatch() → Mailbox.run() 被線程調(diào)度 → Mailbox.processMailbox() → Actor.receive()。遠(yuǎn)程消息傳遞多了序列化、幀編碼、Netty TCP 傳輸、反序列化等步驟。Akka Remote 基于 Netty 實(shí)現(xiàn)分為應(yīng)用層、Remote 層、Netty 層、TCP 層四層。關(guān)鍵特性同一 Actor 的消息順序處理不同 Actor 并行處理無(wú)需鎖和同步。第三Dispatcher 調(diào)度器是 Akka 性能的關(guān)鍵。Dispatcher 負(fù)責(zé)將 Mailbox 調(diào)度到線程池執(zhí)行調(diào)度模型消息 → Mailbox → Dispatcher → ExecutorService → ThreadPool → Actor.receive()。四種調(diào)度器Dispatcher默認(rèn)共享 ForkJoinPool、PinnedDispatcher獨(dú)占線程適合阻塞操作、CallingThreadDispatcher調(diào)用線程僅用于測(cè)試、Custom Dispatcher自定義線程池資源隔離。ForkJoinPool 的工作竊取特性非常適合 Actor 模型。Mailbox 實(shí)現(xiàn) Runnable通過(guò) setAsScheduled() 避免重復(fù)調(diào)度throughput 控制單次處理消息數(shù)。第四配置參數(shù)與性能調(diào)優(yōu)是生產(chǎn)環(huán)境的關(guān)鍵。核心參數(shù)parallelism-factor/min/max線程池大小、throughput單次處理消息數(shù)、mailbox-capacity郵箱容量、heartbeat-interval/acceptable-heartbeat-pause心跳檢測(cè)。調(diào)優(yōu)建議CPU 密集型 parallelism-factor1.0IO 密集型用 PinnedDispatcher 隔離高吞吐調(diào)大 throughput低延遲 throughput1生產(chǎn)環(huán)境用有界郵箱防 OOM避免在 Actor 中執(zhí)行阻塞操作。第五Flink 中 Akka 的使用與優(yōu)化Flink 使用 Akka 作為 RPC 框架JobManager 和 TaskManager 各有獨(dú)立的 ActorSystem。優(yōu)化建議包括線程池配置、心跳檢測(cè)調(diào)優(yōu)、郵箱優(yōu)化、序列化優(yōu)化、避免阻塞操作等。Flink 2.0 已移除 Akka 依賴(lài)改用自研 RPC 框架但 Actor 模型的設(shè)計(jì)思想仍然值得深入學(xué)習(xí)。Akka 的底層設(shè)計(jì)體現(xiàn)了分布式系統(tǒng)的經(jīng)典設(shè)計(jì)思想位置透明、消息驅(qū)動(dòng)、監(jiān)督容錯(cuò)、工作竊取。理解這些底層原理不僅能幫助排查 Flink RPC 問(wèn)題更能體會(huì)到分布式系統(tǒng)設(shè)計(jì)的精妙之處。