@ -26,9 +26,8 @@ import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory ;
import org.apache.curator.framework.imps.CuratorFrameworkState ;
import org.apache.curator.framework.recipes.cache.ChildData ;
import org.apache.curator.framework.recipes.cache.PathChildrenCache ;
import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent ;
import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener ;
import org.apache.curator.framework.recipes.cache.CuratorCache ;
import org.apache.curator.framework.recipes.cache.CuratorCacheListener ;
import org.apache.curator.framework.state.ConnectionState ;
import org.apache.curator.framework.state.ConnectionStateListener ;
import org.apache.curator.retry.RetryForever ;
@ -54,12 +53,12 @@ import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit ;
import java.util.stream.Collectors ;
import static org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent.Type.CHILD_REMOVED ;
import static org.apache.curator.framework.recipes.cache.CuratorCacheAccessor.parentPathFilter ;
@Service
@ConditionalOnProperty ( prefix = "zk" , value = "enabled" , havingValue = "true" , matchIfMissing = false )
@Slf4j
public class ZkDiscoveryService implements DiscoveryService , PathChildrenCacheListener {
public class ZkDiscoveryService implements DiscoveryService {
@Value ( "${zk.url}" )
private String zkUrl ;
@ -84,7 +83,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
private ScheduledExecutorService zkExecutorService ;
@Getter
private CuratorFramework client ;
private PathChildren Cache cache ;
private Curator Cache cache ;
private String nodePath ;
private String zkNodesDir ;
@ -117,7 +116,8 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
@Override
public List < TransportProtos . ServiceInfo > getOtherServers ( ) {
return cache . getCurrentData ( ) . stream ( )
return cache . stream ( )
. filter ( parentPathFilter ( zkNodesDir ) )
. filter ( cd - > ! cd . getPath ( ) . equals ( nodePath ) )
. map ( cd - > {
try {
@ -241,7 +241,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
client = CuratorFrameworkFactory . newClient ( zkUrl , zkSessionTimeout , zkConnectionTimeout , new RetryForever ( zkRetryInterval ) ) ;
client . start ( ) ;
client . blockUntilConnected ( ) ;
cache = new PathChildrenCache ( client , zkNodesDir , true ) ;
cache = CuratorCache . builder ( client , zkNodesDir ) . build ( ) ;
cache . start ( ) ;
stopped = false ;
log . info ( "ZK client connected" ) ;
@ -254,7 +254,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
}
private void subscribeToEvents ( ) {
cache . getL istenable( ) . addListener ( this ) ;
cache . l istenable( ) . addListener ( this : : onCacheEvent ) ;
}
private void unpublishCurrentServer ( ) {
@ -286,29 +286,32 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
return "The " + propertyName + " property need to be set!" ;
}
@Override
public void childEvent ( CuratorFramework curatorFramework , PathChildrenCacheEvent pathChildrenCacheEvent ) throws Exception {
@SneakyThrows
void onCacheEvent ( CuratorCacheListener . Type type , ChildData oldData , ChildData newData ) {
if ( stopped ) {
log . debug ( "Ignoring {}. Service is stopped." , pathChildrenCacheEvent ) ;
log . debug ( "Ignoring {}. Service is stopped." , type ) ;
return ;
}
if ( client . getState ( ) ! = CuratorFrameworkState . STARTED ) {
log . debug ( "Ignoring {}, ZK client is not started, ZK client state [{}]" , pathChildrenCacheEvent , client . getState ( ) ) ;
log . debug ( "Ignoring {}, ZK client is not started, ZK client state [{}]" , type , client . getState ( ) ) ;
return ;
}
ChildData data = pathChildrenCacheEvent . getData ( ) ;
ChildData data = type = = CuratorCacheListener . Type . NODE_DELETED ? oldData : newData ;
if ( data = = null ) {
log . debug ( "Ignoring {} due to empty child data" , pathChildrenCacheEvent ) ;
log . debug ( "Ignoring {} due to empty child data" , type ) ;
return ;
} else if ( data . getData ( ) = = null ) {
log . debug ( "Ignoring {} due to empty child's data" , pathChildrenCacheEvent ) ;
log . debug ( "Ignoring {} due to empty child's data" , type ) ;
return ;
} else if ( zkNodesDir . equals ( data . getPath ( ) ) ) {
log . debug ( "Ignoring event about parent node {}" , data . getPath ( ) ) ;
return ;
} else if ( nodePath ! = null & & nodePath . equals ( data . getPath ( ) ) ) {
if ( pathChildrenCacheEvent . getType ( ) = = CHILD_REMOVED ) {
if ( type = = CuratorCacheListener . Type . NODE_DELET ED) {
log . info ( "ZK node for current instance is somehow deleted." ) ;
publishCurrentServer ( ) ;
}
log . debug ( "Ignoring event about current server {}" , pathChildrenCacheEvent ) ;
log . debug ( "Ignoring event about current server {}" , data . getPath ( ) ) ;
return ;
}
TransportProtos . ServiceInfo instance ;
@ -322,9 +325,9 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
String serviceId = instance . getServiceId ( ) ;
ProtocolStringList serviceTypesList = instance . getServiceTypesList ( ) ;
log . trace ( "Processing [{}] event for [{}]" , pa thChildrenCacheEvent . getT ype( ) , serviceId ) ;
switch ( pa thChildrenCacheEvent . getT ype( ) ) {
case CHILD_ADD ED:
log . trace ( "Processing [{}] event for [{}]" , type , serviceId ) ;
switch ( type ) {
case NODE_CREAT ED:
ScheduledFuture < ? > task = delayedTasks . remove ( serviceId ) ;
if ( task ! = null ) {
if ( task . cancel ( false ) ) {
@ -341,7 +344,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
recalculatePartitions ( ) ;
}
break ;
case CHILD_REMOV ED:
case NODE_DELET ED:
zkExecutorService . submit ( ( ) - > applicationEventPublisher . publishEvent ( new OtherServiceShutdownEvent ( this , serviceId , serviceTypesList ) ) ) ;
ScheduledFuture < ? > future = zkExecutorService . schedule ( ( ) - > {
log . debug ( "[{}] Going to recalculate partitions due to removed node [{}]" ,