You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 01:45:19