Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
package com.github.wenweihu86.raft;

public class PeerId {
Copy link
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

PeerId作为map的key时,应该实现equals和hashcode函数。

private Integer peerId;

public PeerId(Integer peerId) {
this.peerId = peerId;
}

public Integer getPeerId() {
return peerId;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public enum NodeState {

private RaftOptions raftOptions;
private RaftMessage.Configuration configuration;
private ConcurrentMap<Integer, Peer> peerMap = new ConcurrentHashMap<>();
private ConcurrentMap<PeerId, Peer> peerMap = new ConcurrentHashMap<>();
private RaftMessage.Server localServer;
private StateMachine stateMachine;
private SegmentedLog raftLog;
Expand Down Expand Up @@ -131,7 +131,7 @@ public void init() {
&& server.getServerId() != localServer.getServerId()) {
Peer peer = new Peer(server);
peer.setNextIndex(raftLog.getLastLogIndex() + 1);
peerMap.put(server.getServerId(), peer);
peerMap.put(new PeerId(server.getServerId()), peer);
}
}

Expand Down Expand Up @@ -421,7 +421,7 @@ public void applyConfiguration(RaftMessage.LogEntry entry) {
&& server.getServerId() != localServer.getServerId()) {
Peer peer = new Peer(server);
peer.setNextIndex(raftLog.getLastLogIndex() + 1);
peerMap.put(server.getServerId(), peer);
peerMap.put(new PeerId(server.getServerId()), peer);
}
}
LOG.info("new conf is {}, leaderId={}", PRINTER.print(newConfiguration), leaderId);
Expand Down Expand Up @@ -941,6 +941,18 @@ private RaftMessage.InstallSnapshotRequest buildInstallSnapshotRequest(
return requestBuilder.build();
}

public void removePeer(PeerId peerId) {
peerMap.remove(peerId);
}

public boolean containsPeer(PeerId peerId) {
return peerMap.containsKey(peerId);
}

public void addPeer(Peer peer) {
peerMap.put(new PeerId(peer.getServer().getServerId()), peer);
}

public Lock getLock() {
return lock;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package com.github.wenweihu86.raft.service.impl;

import com.github.wenweihu86.raft.Peer;
import com.github.wenweihu86.raft.PeerId;
import com.github.wenweihu86.raft.RaftNode;
import com.github.wenweihu86.raft.proto.RaftMessage;
import com.github.wenweihu86.raft.service.RaftClientService;
Expand Down Expand Up @@ -99,7 +100,7 @@ public RaftMessage.AddPeersResponse addPeers(RaftMessage.AddPeersRequest request
return responseBuilder.build();
}
for (RaftMessage.Server server : request.getServersList()) {
if (raftNode.getPeerMap().containsKey(server.getServerId())) {
if (raftNode.containsPeer(new PeerId(server.getServerId()))) {
LOG.warn("already be added/adding to configuration");
responseBuilder.setResMsg("already be added/adding to configuration");
return responseBuilder.build();
Expand All @@ -110,7 +111,7 @@ public RaftMessage.AddPeersResponse addPeers(RaftMessage.AddPeersRequest request
final Peer peer = new Peer(server);
peer.setNextIndex(1);
requestPeers.add(peer);
raftNode.getPeerMap().putIfAbsent(server.getServerId(), peer);
raftNode.addPeer(peer);
raftNode.getExecutorService().submit(new Runnable() {
@Override
public void run() {
Expand Down Expand Up @@ -163,7 +164,7 @@ public void run() {
try {
for (Peer peer : requestPeers) {
peer.getRpcClient().stop();
raftNode.getPeerMap().remove(peer.getServer().getServerId());
raftNode.removePeer(new PeerId(peer.getServer().getServerId()));
}
} finally {
raftNode.getLock().unlock();
Expand Down