sashapolo commented on code in PR #1366:
URL: https://github.com/apache/ignite-3/pull/1366#discussion_r1032583284


##########
modules/table/src/main/java/org/apache/ignite/internal/table/distributed/TableManager.java:
##########
@@ -1759,95 +1741,95 @@ public boolean onUpdate(@NotNull WatchEvent evt) {
 
                     TablePartitionId replicaGrpId = new 
TablePartitionId(tblId, partId);
 
-                    // Assignments of the pending rebalance that we received 
through the meta storage watch mechanism.
-                    Set<ClusterNode> newPeers = 
ByteUtils.fromBytes(pendingAssignmentsWatchEvent.value());
-
-                    var pendingAssignments = 
metaStorageMgr.get(pendingPartAssignmentsKey(replicaGrpId)).join();
+                    Entry pendingAssignmentsEntry = 
metaStorageMgr.get(pendingPartAssignmentsKey(replicaGrpId)).join();
 
-                    assert pendingAssignmentsWatchEvent.revision() <= 
pendingAssignments.revision()
+                    assert pendingAssignmentsWatchEvent.revision() <= 
pendingAssignmentsEntry.revision()
                             : "Meta Storage watch cannot notify about an event 
with the revision that is more than the actual revision.";
 
+                    // Assignments of the pending rebalance that we received 
through the meta storage watch mechanism.
+                    Set<Assignment> pendingAssignments = 
ByteUtils.fromBytes(pendingAssignmentsWatchEvent.value());
+
                     TableImpl tbl = tablesByIdVv.latest().get(tblId);
 
                     ExtendedTableConfiguration tblCfg = 
(ExtendedTableConfiguration) tablesCfg.tables().get(tbl.name());
 
                     // Stable assignments from the meta store, which revision 
is bounded by the current pending event.
-                    byte[] stableAssignments = 
metaStorageMgr.get(stablePartAssignmentsKey(replicaGrpId),
+                    byte[] stableAssignmentsBytes = 
metaStorageMgr.get(stablePartAssignmentsKey(replicaGrpId),
                             
pendingAssignmentsWatchEvent.revision()).join().value();
 
-                    Set<ClusterNode> assignments = stableAssignments == null
+                    Set<Assignment> stableAssignments = stableAssignmentsBytes 
== null
                             // This is for the case when the first rebalance 
occurs.
-                            ? ((List<Set<ClusterNode>>) 
ByteUtils.fromBytes(tblCfg.assignments().value())).get(partId)
-                            : ByteUtils.fromBytes(stableAssignments);
+                            ? ((List<Set<Assignment>>) 
ByteUtils.fromBytes(tblCfg.assignments().value())).get(partId)
+                            : ByteUtils.fromBytes(stableAssignmentsBytes);
+
+                    List<String> stablePeers = new ArrayList<>();
+                    List<String> stableLearners = new ArrayList<>();
+
+                    for (Assignment assignment : stableAssignments) {
+                        if (assignment.isPeer()) {
+                            stablePeers.add(assignment.consistentId());
+                        } else {
+                            stableLearners.add(assignment.consistentId());
+                        }
+                    }
 
-                    placementDriver.updateAssignment(replicaGrpId, 
assignments);
+                    placementDriver.updateAssignment(replicaGrpId, 
stablePeers);
 
                     ClusterNode localMember = 
raftMgr.topologyService().localMember();
 
-                    List<ClusterNode> deltaPeers = newPeers.stream()
-                            .filter(p -> !assignments.contains(p))
-                            .collect(toList());
+                    // Start a new Raft node and Replica if this node has 
appeared in the new assignments.
+                    boolean shouldStartLocalServices = 
pendingAssignments.stream()
+                            .filter(assignment -> 
!stableAssignments.contains(assignment))
+                            .anyMatch(assignment -> 
localMember.name().equals(assignment.consistentId()));

Review Comment:
   Looks like we don't need to call `filter` at all



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to