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 코디네이터의 동기화 메서드
DistroProtocol의 sync 메서드가 호출됩니다. 여기서는 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메서드는taskExecutionEngineHolder의DelayedTaskEngine에DelayedDistroTask를 추가합니다. 이DelayedTaskEngine은 서비스 변경 알림에서 사용되는PushDelayTaskExecuteEngine과 유사한 역할을 하며,DelayedDistroTask는PushDelayTask에 해당합니다.
DistroTaskEngineHolder는 생성자에서 DelayedTaskEngine의 기본 작업 처리기를 설정합니다.
public DistroTaskEngineHolder(DistroComponentRegistry componentRegistry) {
DistroTaskProcessor defaultProcessor = new DistroTaskProcessor(this, componentRegistry);
delayedTaskEngine.setDefaultTaskProcessor(defaultProcessor);
}
DistroTaskEngineHolder는delayedTaskEngine의 기본 처리기를DistroTaskProcessor로 설정합니다.
Distro 지연 작업 실행 엔진 (DelayedTaskEngine)
DelayedTaskEngine은 PushDelayTaskExecuteEngine과 마찬가지로 NacosDelayTaskExecuteEngine을 상속받습니다.
DelayedTaskEngine은 PushDelayTaskExecuteEngine과 유사하게 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를 생성하여DistroTaskEngineHolder의DistroExecuteTaskExecuteEngine에 추가됩니다.
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()를 호출하여 콜백 메커니즘을 사용하는지 확인합니다. 콜백 메커니즘을 사용하면 작업 완료 후DistroOperationCallback의onSuccess또는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메서드는DataTransportAgent의syncPayload를 호출하여 대상 서버에 데이터를 동기화합니다.
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를 가져옵니다. 여기서ResourceType는Nacos: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 이벤트를 발행합니다. 이 이벤트는 DistroClientDataProcessor와 ClientServiceIndexManager에 의해 처리됩니다. 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());
}
}