我正在尝试实现相同的要求:我有一个需要选举领导者的 Java 服务,并且我没有在 Consul 中配置服务健康检查。
使用来自 Consul-client 的 LeaderElectionUtil 是有问题的,因为如果上面提到的所有原因。不幸的是,也无法自定义 LeaderElectionUtil,因为它的所有内部工作都是使用私有方法完成的(例如,它应该使用 protected 并让用户覆盖会话创建)。
我已尝试按照 consul-client README 中的“基本用法 - 示例 1”中的说明实施“服务注册”,但对我来说是 calling AgentClient.pass() always throws an exception。
所以我的解决方案正是您指定的 - 拥有一个带有 TTL 的会话并在服务存在时更新它。
这是我的实现,它要求用户还注册一个回调,用于检查服务是否仍然有效续订,以防万一:
public class SessionHolder implements Runnable {
private static final String TTL_TEMPLATE = "%ss";
private Consul client;
private String id;
private LinkedList<Supplier<Boolean>> liveChecks = new LinkedList<>();
private long ttl;
private boolean shutdown = false;
public SessionHolder(Consul client, String service, long ttl) {
this.client = client;
this.ttl = ttl;
final Session session = ImmutableSession.builder()
.name(service)
.ttl(String.format(TTL_TEMPLATE, ttl))
.build();
id = client.sessionClient().createSession(session).getId();
Thread upkeep = new Thread(this);
upkeep.setDaemon(true);
upkeep.start();
}
public String getId() {
return id;
}
public void registerKeepAlive(Supplier<Boolean> liveCheck) {
liveChecks.add(liveCheck);
}
@Override
public synchronized void run() {
// don't start renewing immediately
try {
wait(ttl / 2 * 1000);
} catch (InterruptedException e) {}
while (!isShutdown()) {
if (liveChecks.isEmpty() || liveChecks.stream().allMatch(Supplier::get)) {
client.sessionClient().renewSession(getId());
}
try {
wait(ttl / 2 * 1000);
} catch (InterruptedException e) {
// go on, try again
}
}
}
public synchronized boolean isShutdown() {
return shutdown;
}
public synchronized void close() {
shutdown = true;
notify();
client.sessionClient().destroySession(getId());
}
}
那么选举领导者或多或少就像这样简单:
if (consul.keyValueClient().acquireLock(getServiceKey(service), currentNode, sessionHolder.getId()))
return true; // I'm the leader
需要记住的一件事是,如果会话在没有正确清理的情况下终止(我在上面的 SessionHolder.close() 中所做的),consul 的 lock-delay 功能将阻止新领导者被选举约 15 秒(默认值,不幸的是 Consul-client 不提供 API 来修改)。
为了解决这个问题,除了确保如上所示正确终止服务后自行清理之外,我还确保让服务在所需的最少时间内保持领导位置,并在以下情况下释放领导不再使用它,请致电consul.keyValueClient().releaseLock()。例如,我有一个集群服务,我们选择一个领导者从外部 RDBMS 读取数据更新(然后直接分布在集群中,而不是每个节点重新加载所有数据)。由于这是通过轮询完成的,每个节点都会在轮询之前尝试被选中,如果被选中,它将轮询数据库,传播更新并退出。如果之后崩溃,delay-lock 将不会阻止另一个节点轮询。