Nacos 서비스 인스턴스 데이터의 클러스터 동기화 과정

Nacos 서버가 서비스 등록 요청을 처리할 때, 내부적으로 클라이언트 변경 이벤트(ClientChangedEvent)를 발생시킵니다. 이 이벤트는 주로 임시 클라이언트 작업 서비스 구현체(EphemeralClientOperationServiceImpl) 내의 addServiceInstance 메서드 호출 체인을 통해 시작됩니다.

public boolean addServiceEntry(ServiceIdentity svcIdentifier, InstanceRegistrationInfo registrationPayload) {
    // 최초 등록 시, 해당 서비스에 대한 이전 정보가 없으므로 null 반환
    if (null == activePublishers.put(svcIdentifier, registrationPayload)) {
        // 메트릭 카운트 증가
        Metrics.incrementActiveInstanceCount();
    }
    // 클라이언트 상태 변경 이벤트 발행
    EventPublisher.publishEvent(new ClientStateEvent.ClientChangedEvent(this));
    ServiceLogger.info("Client data changed for service {}, {}", svcIdentifier, getClientIdentifier());
    return true;
}

클라이언트 변경 이벤트 처리

ClientChangedEvent는 분산 클라이언트 데이터 처리기(DistroClientDataProcessor)의 onClientEvent 메서드에서 처리됩니다.

@Override
public void onClientEvent(ApplicationEvent event) {
    // 독립 실행 모드(Standalone Mode)에서는 분산 처리 불필요
    if (ServerConfigUtil.isSingleNodeMode()) {
        return;
    }
    // 이벤트 유형에 따른 처리 분기
    if (event instanceof ClientStateEvent.ClientVerificationFailedEvent) {
        synchronizeFailedVerificationToNode((ClientStateEvent.ClientVerificationFailedEvent) event);
    } else {
        // 일반 클라이언트 변경 이벤트 처리
        synchronizeClientUpdatesToAllNodes((ClientStateEvent) event);
    }
}

private void synchronizeClientUpdatesToAllNodes(ClientStateEvent event) {
    ClientContext currentClient = event.getClient();
    // 클라이언트가 임시(ephemeral) 데이터이며, 현재 노드가 담당하는 클라이언트인 경우에만 동기화
    if (currentClient == null || !currentClient.isEphemeral() || !nodeManager.isResponsibleForClient(currentClient)) {
        return;
    }
    // 클라이언트 연결 해제 이벤트인 경우 삭제 작업으로 동기화
    if (event instanceof ClientStateEvent.ClientDisconnectedEvent) {
        DistroResourceKey resourceId = new DistroResourceKey(currentClient.getClientIdentifier(), RESOURCE_TYPE);
        distroCoordinator.syncResource(resourceId, DistroOperation.REMOVE);
    } else if (event instanceof ClientStateEvent.ClientChangedEvent) {
        // 클라이언트 변경 이벤트인 경우 변경 작업으로 동기화
        DistroResourceKey resourceId = new DistroResourceKey(currentClient.getClientIdentifier(), RESOURCE_TYPE);
        distroCoordinator.syncResource(resourceId, DistroOperation.UPDATE);
    }
}

Distro 코디네이터의 동기화 메서드

DistroProtocolsync 메서드가 호출됩니다. 여기서는 syncResource로 이름이 변경되었습니다.

public void syncResource(DistroResourceKey key, DistroOperation operation) {
    syncResource(key, operation, DistroConfiguration.getInstance().getDefaultSyncDelayMs());
}

public void syncResource(DistroResourceKey key, DistroOperation operation, long delayDuration) {
    for (ClusterNode node : clusterNodeManager.getAllRemoteNodes()) {
        initiateSyncToTargetNode(key, operation, node.getNodeAddress(), delayDuration);
    }
}

public void initiateSyncToTargetNode(DistroResourceKey key, DistroOperation operation, String targetNodeAddress, long delayDuration) {
    DistroResourceKey keyWithTarget = new DistroResourceKey(key.getResourceIdentifier(), key.getResourceType(), targetNodeAddress);
    DelayedDistroTask delayedTask = new DelayedDistroTask(keyWithTarget, operation, delayDuration);
    taskExecutionEngineHolder.getDelayedTaskEngine().submitTask(keyWithTarget, delayedTask);

    if (DistroLogger.isDebugEnabled()) {
        DistroLogger.debug("[DISTRO-SCHEDULE] {} to {}", key, targetNodeAddress);
    }
}
  • initiateSyncToTargetNode 메서드는 taskExecutionEngineHolderDelayedTaskEngineDelayedDistroTask를 추가합니다. 이 DelayedTaskEngine은 서비스 변경 알림에서 사용되는 PushDelayTaskExecuteEngine과 유사한 역할을 하며, DelayedDistroTaskPushDelayTask에 해당합니다.

DistroTaskEngineHolder는 생성자에서 DelayedTaskEngine의 기본 작업 처리기를 설정합니다.

public DistroTaskEngineHolder(DistroComponentRegistry componentRegistry) {
    DistroTaskProcessor defaultProcessor = new DistroTaskProcessor(this, componentRegistry);
    delayedTaskEngine.setDefaultTaskProcessor(defaultProcessor);
}
  • DistroTaskEngineHolderdelayedTaskEngine의 기본 처리기를 DistroTaskProcessor로 설정합니다.

Distro 지연 작업 실행 엔진 (DelayedTaskEngine)

DelayedTaskEnginePushDelayTaskExecuteEngine과 마찬가지로 NacosDelayTaskExecuteEngine을 상속받습니다.

DelayedTaskEnginePushDelayTaskExecuteEngine과 유사하게 DelayedDistroTask를 처리하기 위해 DistroTaskProcessor를 스케줄링합니다.

// DistroTaskProcessor
@Override
public boolean process(AbstractNacosTask task) {
    if (!(task instanceof DelayedDistroTask)) {
        return false;
    }
    DelayedDistroTask delayedDistroTask = (DelayedDistroTask) task;
    DistroResourceKey resourceKey = delayedDistroTask.getResourceKey();

    switch (delayedDistroTask.getOperation()) {
        case REMOVE:
            // 삭제 작업
            DistroSyncDeletionTask deleteTask = new DistroSyncDeletionTask(resourceKey, componentRegistry);
            taskExecutionEngineHolder.getWorkerPoolManager().submitTask(resourceKey, deleteTask);
            return true;
        case UPDATE:
        case ADD:
            // 변경 또는 추가 작업
            DistroSyncUpdateTask updateTask = new DistroSyncUpdateTask(resourceKey, componentRegistry);
            taskExecutionEngineHolder.getWorkerPoolManager().submitTask(resourceKey, updateTask);
            return true;
        default:
            return false;
    }
}
  • 삭제, 변경, 추가 작업 모두 DistroSyncUpdateTask 또는 DistroSyncDeletionTask를 생성하여 DistroTaskEngineHolderDistroExecuteTaskExecuteEngine에 추가됩니다.

DistroSyncUpdateTask의 실행 메서드 (run)

@Override
public void run() {
    // Nacos:Naming:v2:ClientData
    String dataType = getResourceKey().getResourceType();

    // DistroClientTransportAgent 객체 가져오기
    DataTransportAgent transportAgent = componentRegistry.findTransportAgent(dataType);
    if (null == transportAgent) {
        DistroLogger.warn("No transport agent found for data type [{}]", dataType);
        return;
    }
    DistroLogger.info("[DISTRO-START] {}", toString());

    // 기본적으로 true 반환
    if (transportAgent.supportsCallbackTransport()) {
        // 이 경로를 통해 실행
        executeWithCallbackHandler(new DistroOperationCallback());
    } else {
        performDistroOperation();
    }
}
  • 먼저 resourceKey에서 ResourceType(기본값: Nacos:Naming:v2:ClientData)를 가져옵니다.
  • 해당 타입에 맞는 DataTransportAgent를 찾은 후, supportsCallbackTransport()를 호출하여 콜백 메커니즘을 사용하는지 확인합니다. 콜백 메커니즘을 사용하면 작업 완료 후 DistroOperationCallbackonSuccess 또는 onFailed 메서드가 호출됩니다.

DistroSyncUpdateTask의 executeWithCallbackHandler 메서드

@Override
protected void executeWithCallbackHandler(DistroCallback callback) {
    String dataType = getResourceKey().getResourceType();
    // 요청 데이터 가져오기
    DistroPayload distroPayload = retrieveDistroPayload(dataType);
    if (null == distroPayload) {
        DistroLogger.warn("[DISTRO] {} with null payload to sync, skipping", toString());
        return;
    }
    // syncPayload를 사용하여 클러스터 노드 동기화
    componentRegistry.findTransportAgent(dataType)
            .syncPayload(distroPayload, getResourceKey().getTargetNodeAddress(), callback);
}
  • executeWithCallbackHandler 메서드는 DataTransportAgentsyncPayload를 호출하여 대상 서버에 데이터를 동기화합니다.

DataTransportAgent의 syncPayload 메서드

@Override
public void syncPayload(DistroPayload payload, String targetAddress, DistroCallback callback) {
    if (isTargetNodeNonExistent(targetAddress)) {
        callback.onSuccess();
        return;
    }
    // 요청 객체 생성
    DistroDataRequest request = new DistroDataRequest(payload, payload.getDataType());

    // 클러스터 노드 찾기
    ClusterNode memberNode = clusterNodeManager.findNode(targetAddress);
    try {
        // 비동기 RPC 요청 전송
        clusterRpcGateway.sendAsyncRequest(memberNode, request, new DistroRpcCallbackWrapper(callback, memberNode));
    } catch (NacosException e) {
        callback.onFailed(e);
    }
}
  • syncPayload 메서드는 DistroDataRequest를 생성한 후, 클러스터의 RPC 클라이언트를 사용하여 다른 Nacos 서버에 서비스 정보 동기화 요청을 보냅니다.

DistroDataRequest 처리 로직

DistroDataRequest 요청은 결국 다른 Nacos 서버의 DistroDataRequestHandler.handleRequest 메서드에서 처리됩니다.

// DistroDataRequestHandler
@Override
public DistroDataResponse handleRequest(DistroDataRequest request, RequestMetadata metadata) throws NacosException {
    try {
        switch (request.getOperationType()) {
            case VERIFY:
                return processVerification(request.getPayload(), metadata);
            case SNAPSHOT:
                return generateSnapshot();
            case ADD:
            case UPDATE:
            case REMOVE:
                // 추가, 변경, 삭제 모두 이 메서드를 통해 처리
                return processSyncData(request.getPayload());
            case QUERY:
                return processQueryData(request.getPayload());
            default:
                return new DistroDataResponse();
        }
    } catch (Exception e) {
        DistroLogger.error("[DISTRO-FAILED] Distro request handling failed with exception", e);
        DistroDataResponse response = new DistroDataResponse();
        response.setErrorCode(ErrorCode.FAILURE.getCode());
        response.setMessage("Failed to handle distro request due to exception");
        return response;
    }
}
  • 수신된 동기화 요청이 추가, 변경, 삭제 중 무엇이든 processSyncData() 메서드를 통해 처리됩니다.
// DistroDataRequestHandler
private DistroDataResponse processSyncData(DistroPayload payload) {
    DistroDataResponse response = new DistroDataResponse();
    // onReceive 처리 메서드 호출
    if (!distroCoordinator.onReceive(payload)) {
        response.setErrorCode(ErrorCode.FAILURE.getCode());
        response.setMessage("[DISTRO-FAILED] Distro payload processing failed");
    }
    return response;
}

// DistroCoordinator의 처리 로직
public boolean onReceive(DistroPayload receivedPayload) {
    DistroLogger.info("[DISTRO] Received distro payload type: {}, key: {}", receivedPayload.getDataType(),
            receivedPayload.getResourceKey());
    // ResourceType은 여전히 Nacos:Naming:v2:ClientData
    String resourceType = receivedPayload.getResourceKey().getResourceType();

    // DistroClientTransportAgent 처리 객체 가져오기
    PayloadProcessor processor = componentRegistry.findPayloadProcessor(resourceType);
    if (null == processor) {
        DistroLogger.warn("[DISTRO] No payload processor found for received data {}", resourceType);
        return false;
    }
    // 최종적으로 이 메서드 호출
    return processor.handlePayload(receivedPayload);
}
  • processSyncData 메서드는 먼저 ResourceType를 가져옵니다. 여기서 ResourceTypeNacos:Naming:v2:ClientData에 해당합니다. 해당 타입에 따라 findPayloadProcessor()를 호출하며, 이 안에는 Map이 있어 DistroClientDataProcessor 객체를 가져옵니다. 마지막으로 handlePayload()를 호출합니다.
// DistroClientDataProcessor
@Override
public boolean handlePayload(DistroPayload payload) {
    switch (payload.getOperation()) {
        case ADD:
        case UPDATE:
            // 추가 및 변경은 이 로직을 따릅니다.
            ClientSyncData clientData = ApplicationUtils.getBean(DataSerializer.class)
                    .deserialize(payload.getContent(), ClientSyncData.class);
            processClientSynchronizationData(clientData);
            return true;
        case REMOVE:
            // 삭제는 이 로직을 따릅니다.
            String clientIdToRemove = payload.getResourceKey().getResourceIdentifier();
            DistroLogger.info("[Client-Remove] Received distro client sync data {}", clientIdToRemove);
            clientConnectionManager.disconnectClient(clientIdToRemove);
            return true;
        default:
            return false;
    }
}
  • distroPayload가 추가 또는 변경 유형일 경우 processClientSynchronizationData() 메서드가 호출됩니다.
// DistroClientDataProcessor
private void processClientSynchronizationData(ClientSyncData syncData) {
    DistroLogger.info("[Client-Update] Received distro client sync data {}, revision={}", syncData.getClientIdentifier(),
            syncData.getAttributes().getAttribute(ClientAttributeKeys.REVISION_NUM, 0L));
    clientConnectionManager.updateClientConnection(syncData.getClientIdentifier(), syncData.getAttributes());
    ClientContext existingClient = clientConnectionManager.getClientContext(syncData.getClientIdentifier());
    updateClientState(existingClient, syncData);
}
  • processClientSynchronizationData() 메서드는 다시 updateClientState()를 호출합니다.
private void updateClientState(ClientContext targetClient, ClientSyncData incomingData) {
    // 동기화 데이터에서 네임스페이스 목록 가져오기
    List<String> namespaces = incomingData.getNamespaceList();
    // 동기화 데이터에서 그룹 이름 목록 가져오기
    List<String> groupNames = incomingData.getGroupNamesList();
    // 동기화 데이터에서 서비스 이름 목록 가져오기
    List<String> serviceNames = incomingData.getServiceNamesList();
    // 동기화 데이터에서 인스턴스 게시 정보 목록 가져오기
    List<InstanceRegistrationInfo> instanceInfos = incomingData.getInstanceRegistrationInfos();
    // 동기화된 서비스를 저장할 HashSet 생성
    Set<ServiceIdentity> updatedServices = new HashSet<>();

    // 네임스페이스, 그룹 이름, 서비스 이름 목록을 순회하며 서비스 객체 생성 및 동기화 처리
    for (int i = 0; i < namespaces.size(); i++) {
        // 네임스페이스, 그룹 이름, 서비스 이름으로 서비스 객체 생성
        ServiceIdentity svcKey = ServiceIdentity.create(namespaces.get(i), groupNames.get(i), serviceNames.get(i));
        // 서비스 매니저에서 싱글톤 서비스 객체 가져오기
        ServiceIdentity managedSvc = ServiceRegistry.getInstance().getManagedService(svcKey);
        // 싱글톤 서비스 객체를 업데이트된 서비스 집합에 추가
        updatedServices.add(managedSvc);
        // 현재 인덱스에 해당하는 인스턴스 게시 정보 가져오기
        InstanceRegistrationInfo currentInstanceInfo = instanceInfos.get(i);
        // 현재 인스턴스 게시 정보가 클라이언트의 해당 서비스 인스턴스 게시 정보와 다를 경우
        if (!currentInstanceInfo.equals(targetClient.getInstanceInfoForService(managedSvc))) {
            // 클라이언트에 서비스 인스턴스 추가 및 서비스 등록 이벤트 발행
            targetClient.addServiceInstance(managedSvc, currentInstanceInfo);
            EventPublisher.publishEvent(
                    new ClientLifecycleEvent.ClientRegisteredServiceEvent(managedSvc, targetClient.getClientIdentifier()));
        }
    }

    // 클라이언트가 게시한 모든 서비스를 순회
    for (ServiceIdentity existingSvc : targetClient.getAllPublishedServices()) {
        // 현재 서비스가 업데이트된 서비스 집합에 포함되어 있지 않으면
        if (!updatedServices.contains(existingSvc)) {
            // 클라이언트에서 서비스 인스턴스 제거 및 서비스 등록 해제 이벤트 발행
            targetClient.removeServiceInstance(existingSvc);
            EventPublisher.publishEvent(
                    new ClientLifecycleEvent.ClientDeregisteredServiceEvent(existingSvc, targetClient.getClientIdentifier()));
        }
    }
}
  • updateClientState의 로직은 서버가 서비스 등록 요청을 처리하는 방식과 유사합니다. 주로 클라이언트의 addServiceInstance 메서드를 호출하여 서비스와 인스턴스 관계를 클라이언트에 등록합니다. 이후 ClientOperationEvent.ClientRegisterServiceEvent 이벤트를 발행합니다.

만약 동기화 정보가 REMOVE 유형이라면 clientConnectionManager.disconnectClient 메서드가 호출됩니다.

// EphemeralClientConnectionManager
@Override
public boolean disconnectClient(String clientId) {
    ServiceLogger.info("Client connection {} disconnected, removing instances and subscribers", clientId);
    // 클라이언트 정보 제거
    ClientContext clientToRemove = activeClientContexts.remove(clientId);
    if (null == clientToRemove) {
        return true;
    }
    // 클라이언트 연결 해제 이벤트 발행
    EventPublisher.publishEvent(new ClientStateEvent.ClientDisconnectedEvent(clientToRemove));
    clientToRemove.releaseResources();
    return true;
}

disconnectClient 메서드는 클라이언트 정보를 제거한 후 ClientStateEvent.ClientDisconnectedEvent 이벤트를 발행합니다. 이 이벤트는 DistroClientDataProcessorClientServiceIndexManager에 의해 처리됩니다. DistroClientDataProcessor는 연결 해제 이벤트를 다른 서버에 동기화하는 데 사용되며, ClientServiceIndexManager는 현재 서버에서 클라이언트 연결 해제 관련 로직을 수행합니다.

// ClientServiceIndexManager
@Override
public void onClientEvent(ApplicationEvent event) {
    if (event instanceof ClientStateEvent.ClientDisconnectedEvent) {
        // 클라이언트 연결 해제 처리
        handleClientDisconnection((ClientStateEvent.ClientDisconnectedEvent) event);
    } else if (event instanceof ClientLifecycleEvent) {
        handleClientLifecycleOperation((ClientLifecycleEvent) event);
    }
}

private void handleClientDisconnection(ClientStateEvent.ClientDisconnectedEvent event) {
    ClientContext client = event.getClient();
    for (ServiceIdentity subscribedSvc : client.getAllSubscribedServices()) {
        // 구독자 정보 제거
        removeSubscriberMappings(subscribedSvc, client.getClientIdentifier());
    }
    for (ServiceIdentity publishedSvc : client.getAllPublishedServices()) {
        // 등록 정보 제거
        removePublisherMappings(publishedSvc, client.getClientIdentifier());
    }
}

태그: nacos service-discovery clustering distributed-systems event-driven

8월 18일 23:51에 게시됨