1111// limitations under the License.
1212package edu .iu .dsc .tws .checkpointing .master ;
1313
14- import java .util .HashSet ;
15- import java .util .Set ;
14+ import java .util .HashMap ;
1615import java .util .logging .Logger ;
1716
1817import edu .iu .dsc .tws .api .net .request .RequestID ;
@@ -24,7 +23,7 @@ public class FamilyInitHandler {
2423 private static final Logger LOG = Logger .getLogger (FamilyInitHandler .class .getName ());
2524
2625 private int count ;
27- private Set < RequestID > pendingResponses ;
26+ private HashMap < Integer , RequestID > pendingResponses ;
2827 private RRServer rrServer ;
2928 private String family ;
3029 private Long familyVersion ;
@@ -35,20 +34,25 @@ public FamilyInitHandler(RRServer rrServer,
3534 this .rrServer = rrServer ;
3635 this .family = family ;
3736 this .familyVersion = familyVersion ;
38- this .pendingResponses = new HashSet <>();
37+ this .pendingResponses = new HashMap <>();
3938 this .count = count ;
4039 }
4140
42- public boolean scheduleResponse (RequestID requestID ) {
43- this .pendingResponses .add (requestID );
41+ public boolean scheduleResponse (int workerId , RequestID requestID ) {
42+ RequestID previousRequest = this .pendingResponses .put (workerId , requestID );
43+ if (previousRequest != null ) {
44+ LOG .warning ("Duplicate request received for " + this .family
45+ + " from worker : " + workerId + ". Workers might be coming after a failure." );
46+ }
4447 if (this .pendingResponses .size () == count ) {
45- for (RequestID pendingRespons : this .pendingResponses ) {
48+ for (RequestID pendingRespons : this .pendingResponses . values () ) {
4649 this .rrServer .sendResponse (pendingRespons ,
4750 Checkpoint .FamilyInitializeResponse .newBuilder ()
4851 .setFamily (this .family )
4952 .setVersion (this .familyVersion )
5053 .build ());
5154 }
55+ this .pendingResponses .clear ();
5256 return true ;
5357 } else {
5458 return false ;
0 commit comments