fix issue 6504, and polish some code (#6505)
This commit is contained in:
parent
d41bfaba21
commit
47ed754892
|
|
@ -32,6 +32,7 @@ import com.orbitz.consul.Consul;
|
|||
import com.orbitz.consul.KeyValueClient;
|
||||
import com.orbitz.consul.cache.KVCache;
|
||||
import com.orbitz.consul.model.kv.Value;
|
||||
import org.apache.dubbo.common.utils.StringUtils;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.LinkedHashSet;
|
||||
|
|
@ -55,17 +56,26 @@ public class ConsulDynamicConfiguration extends TreePathDynamicConfiguration {
|
|||
private static final int DEFAULT_WATCH_TIMEOUT = 60 * 1000;
|
||||
private static final String WATCH_TIMEOUT = "consul-watch-timeout";
|
||||
|
||||
private Consul client;
|
||||
private final Consul client;
|
||||
|
||||
private KeyValueClient kvClient;
|
||||
private final KeyValueClient kvClient;
|
||||
|
||||
private ConcurrentMap<String, ConsulListener> watchers = new ConcurrentHashMap<>();
|
||||
private final int watchTimeout;
|
||||
|
||||
private final ConcurrentMap<String, ConsulListener> watchers = new ConcurrentHashMap<>();
|
||||
|
||||
public ConsulDynamicConfiguration(URL url) {
|
||||
super(url);
|
||||
watchTimeout = url.getParameter(WATCH_TIMEOUT, DEFAULT_WATCH_TIMEOUT);
|
||||
String host = url.getHost();
|
||||
int port = url.getPort() != 0 ? url.getPort() : DEFAULT_PORT;
|
||||
client = Consul.builder().withHostAndPort(HostAndPort.fromParts(host, port)).build();
|
||||
Consul.Builder builder = Consul.builder()
|
||||
.withHostAndPort(HostAndPort.fromParts(host, port));
|
||||
String token = url.getParameter("token", (String) null);
|
||||
if (StringUtils.isNotEmpty(token)) {
|
||||
builder.withAclToken(token);
|
||||
}
|
||||
client = builder.build();
|
||||
this.kvClient = client.keyValueClient();
|
||||
}
|
||||
|
||||
|
|
@ -128,8 +138,8 @@ public class ConsulDynamicConfiguration extends TreePathDynamicConfiguration {
|
|||
private class ConsulListener implements KVCache.Listener<String, Value> {
|
||||
|
||||
private KVCache kvCache;
|
||||
private Set<ConfigurationListener> listeners = new LinkedHashSet<>();
|
||||
private String normalizedKey;
|
||||
private final Set<ConfigurationListener> listeners = new LinkedHashSet<>();
|
||||
private final String normalizedKey;
|
||||
|
||||
public ConsulListener(String normalizedKey) {
|
||||
this.normalizedKey = normalizedKey;
|
||||
|
|
@ -137,7 +147,7 @@ public class ConsulDynamicConfiguration extends TreePathDynamicConfiguration {
|
|||
}
|
||||
|
||||
private void initKVCache() {
|
||||
this.kvCache = KVCache.newCache(kvClient, normalizedKey);
|
||||
this.kvCache = KVCache.newCache(kvClient, normalizedKey, watchTimeout);
|
||||
kvCache.addListener(this);
|
||||
kvCache.start();
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue