Akka Typed ActorSystem种子节点无法加入集群求助
Akka Typed集群种子节点无法加入问题排查与解决
问题概述
搭建Akka Typed ActorSystem集群时,种子节点无法加入,出现以下错误:
Cluster Node [akka://KYCSystem@127.0.1.1:25520] - Joining of seed-nodes [akka://ClusterSystem@127.0.0.1:25520] was unsuccessful after configured shutdown-after-unsuccessful-join-seed-nodes [20000 milliseconds]. Running CoordinatedShutdown.
相关代码与配置
ActorSystem实现代码
package io.tajji.kycpipeline.actor import akka.actor.typed.ActorSystem import akka.actor.typed.javadsl.Behaviors import akka.cluster.typed.Cluster import io.tajji.apis.kyc.events.* import io.tajji.kycpipeline.model.Command.* import io.tajji.kycpipeline.model.Message import io.tajji.kycpipeline.model.Message.* import io.tajji.kycpipeline.service.KYCValidationService import io.tajji.kycpipeline.service.TextExtractionService import io.tajji.kycpipeline.service.TextParsingService import org.axonframework.eventhandling.EventBus import org.axonframework.eventhandling.GenericEventMessage import org.springframework.stereotype.Component import reactor.core.publisher.Mono @Component class KYCSystem( private val eventBus: EventBus, private val textExtractor: TextExtractionService, private val detailsParser: TextParsingService, private val kycValidator: KYCValidationService ) { private val system: ActorSystem<Message> = ActorSystem.create( Behaviors.setup<Message> { context -> val validator = context.spawn(KYCValidator.create(kycValidator), "DocumentValidatorActor") val documentExtractorActor = context.spawn(DocumentExtractor.create(textExtractor), "TextExtractorActor") val detailsParserActor = context.spawn(DocumentParser.create(detailsParser), "DocumentParserActor") Behaviors.receiveMessage { message -> when(message) { is LandlordKYC -> { val validateLandlordDocuments = ValidateLandlordDocuments( message.event.accountId, message.event, context.self ) validator.tell(validateLandlordDocuments) } is ResidentKYC -> { val validateResidentDocuments = ValidateResidentDocuments( message.event.accountId, message.event, context.self ) validator.tell(validateResidentDocuments) } is ValidLandlordDocuments -> { val extractLandlordDetails = ExtractLandlordDetails( message.accountId, message.idFront, message.idBack, message.pinCertificateData, context.self ) documentExtractorActor.tell(extractLandlordDetails) } is ValidResidentDocuments -> { val extractResidentDetails = ExtractResidentDetails( message.accountId, message.idFront, message.idBack, context.self ) documentExtractorActor.tell(extractResidentDetails) } is ExtractedLandlordDetails -> { val parseLandlordDetails = ParseLandlordDetails( message.accountId, message, context.self ) detailsParserActor.tell(parseLandlordDetails) } is ExtractedResidentDetails -> { val parseResidentDetails = ParseResidentDetails( message.accountId, message, context.self ) detailsParserActor.tell(parseResidentDetails) } is NationalIdInvalid -> { eventBus.publish(GenericEventMessage .asEventMessage<KYCFailed>(KYCFailed( message.accountId, message.reason ))) } is PinCertificateInvalid -> { eventBus.publish(GenericEventMessage .asEventMessage<KYCFailed>(KYCFailed( message.accountId, message.reason ))) } is ParsedResidentDetails -> { message.nationalIDData.subscribe { data -> eventBus.publish(GenericEventMessage .asEventMessage<ResidentKYCPassed>(ResidentKYCPassed( message.accountId, data ))) } } is ParsedLandlordDetails -> { val accountId = message.accountId Mono.zip( message.nationalIDData, message.taxData ).subscribe { tuple -> val nationalId = tuple.t1 val taxData = tuple.t2 eventBus.publish(GenericEventMessage .asEventMessage<LandlordKYCPassed>(LandlordKYCPassed( accountId, nationalId, taxData ))) } } is Invalid -> { eventBus.publish(GenericEventMessage .asEventMessage<KYCFailed>(KYCFailed( message.accountId, message.failureReason ))) } } Behaviors.same() } }, "KYCSystem" ) private val kycCluster = Cluster.get(system) fun processLandlordKYC(event: LandlordKYCRequested) { system.tell(LandlordKYC(event)) } fun processResidentKYC(event: ResidentKYCRequested) { system.tell(ResidentKYC(event)) } }
集群配置文件
akka { cluster { seed-nodes = ["akka://ClusterSystem@127.0.0.1:25520"] shutdown-after-unsuccessful-join-seed-nodes = 20s seed-node-timeout = 15s log-info-verbose = off downing-provider-class = "akka.cluster.sbr.SplitBrainResolverProvider" } actor { provider = cluster } remote { netty.tcp { hostname = "127.0.0.1" port = 25520 } artery { } } coordinated-shutdown { exit-jvm = on } }
排查与修复步骤
1. 修正ActorSystem名称不匹配问题
错误日志显示本地节点名称为KYCSystem,但配置中种子节点名称是ClusterSystem,两者必须完全一致。
- 方案一:修改ActorSystem创建时的名称为
ClusterSystem:private val system: ActorSystem<Message> = ActorSystem.create(..., "ClusterSystem") - 方案二:修改配置文件中的seed-nodes为
akka://KYCSystem@127.0.0.1:25520
2. 统一网络地址
本地节点绑定的是127.0.1.1,但配置中指定的是127.0.0.1,这是系统hosts解析导致的差异。
- 在配置文件中显式指定正确的hostname,或者使用
0.0.0.0绑定所有可用地址:remote.artery.canonical.hostname = "127.0.0.1"
3. 清理远程通信配置冲突
Akka Typed推荐使用Artery作为远程通信层,配置中同时存在netty.tcp和artery会导致冲突,移除netty.tcp配置,改用Artery配置:
remote { artery { canonical.hostname = "127.0.0.1" canonical.port = 25520 } }
4. 显式触发集群加入操作
当前代码仅获取了Cluster实例,但未执行加入集群的逻辑,需添加初始化代码:
private val kycCluster = Cluster.get(system) init { // 如果当前节点是种子节点,调用joinSelf() // kycCluster.joinSelf() // 如果加入已启动的种子节点,使用join方法 kycCluster.join(akka.actor.Address("akka", "ClusterSystem", "127.0.0.1", 25520)) }
5. 确保种子节点先启动
启动节点时,必须先启动配置中的种子节点,待其完全启动后,再启动其他节点尝试加入。
内容的提问来源于stack exchange,提问作者Arnold Odhiambo
相关产品推荐
相关产品推荐

