Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
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
Expand Up @@ -206,7 +206,7 @@ private boolean startBasicService() {
return false;
}
}

// 同步controller的元数据
schedulingSyncBrokerMetadata();

// Register syncStateSet changed listener.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -116,9 +116,15 @@ private void initPeerIdMap() {
String[] peers = this.controllerConfig.getJraftConfig().getjRaftInitConf().split(",");
String[] rpcAddrs = this.controllerConfig.getJraftConfig().getjRaftControllerRPCAddr().split(",");
for (int i = 0; i < peers.length; i++) {
String peerStr = peers[i];
// Remove /learner suffix if present, as PeerId.parse() cannot handle it
int learnerIndex;
if ((learnerIndex = peerStr.indexOf("/learner")) > 0) {
peerStr = peerStr.substring(0, learnerIndex);
}
PeerId peerId = new PeerId();
if (!peerId.parse(peers[i])) {
throw new IllegalArgumentException("Fail to parse peerId:" + peers[i]);
if (!peerId.parse(peerStr)) {
throw new IllegalArgumentException("Fail to parse peerId:" + peerStr);
}
this.peerIdToAddr.put(peerId, rpcAddrs[i]);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ public class BrokerReplicaInfo implements Serializable {

// Start from 1
private final AtomicLong nextAssignBrokerId;

//registerCheckCode值:this.brokerAddress + ";" + System.currentTimeMillis(),作用后续ApplyBrokerId动作的身份校验
private final Map<Long/*brokerId*/, Pair<String/*ipAddress*/, String/*registerCheckCode*/>> brokerIdInfo;

public BrokerReplicaInfo(String clusterName, String brokerName) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,9 +28,11 @@
public class SyncStateInfo implements Serializable {
private final String clusterName;
private final String brokerName;
//master版本号(任期号)
private final AtomicInteger masterEpoch;
//同步副本数集合的版本(任期号)
private final AtomicInteger syncStateSetEpoch;

//同步副本数集合包含master
private Set<Long/*brokerId*/> syncStateSet;

private Long masterBrokerId;
Expand Down
2 changes: 1 addition & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,7 @@
<opentelemetry-exporter-prometheus.version>1.29.0-alpha</opentelemetry-exporter-prometheus.version>
<jul-to-slf4j.version>2.0.6</jul-to-slf4j.version>
<s3.version>2.20.29</s3.version>
<rocksdb.version>1.0.2</rocksdb.version>
<rocksdb.version>1.0.3</rocksdb.version>
<jackson-databind.version>2.13.4.2</jackson-databind.version>
<sofa-jraft.version>1.3.14</sofa-jraft.version>

Expand Down