Skip to content

Commit aa6862e

Browse files
committed
Merge branch 'master' into ahmet/fault-tolerance
2 parents 01fb7f8 + 8a21c54 commit aa6862e

12 files changed

Lines changed: 81 additions & 34 deletions

File tree

‎docs/website/pages/en/configs.js‎

Lines changed: 12 additions & 18 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎twister2/api/src/java/edu/iu/dsc/tws/api/data/FileSystem.java‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,8 +133,23 @@ public abstract BlockLocation[] getFileBlockLocations(FileStatus file,
133133

134134
public abstract FSDataInputStream open(final Path f, final int bufferSize) throws IOException;
135135

136+
/**
137+
* Create a file system with the OVERWRITE
138+
* @param f path
139+
* @return the output stream
140+
* @throws IOException if an error occurs
141+
*/
136142
public abstract FSDataOutputStream create(final Path f) throws IOException;
137143

144+
/**
145+
* Create a file system with the specific write mdoe
146+
* @param f path
147+
* @param writeMode weather overwrite or not, when creating new files
148+
* @return the output stream
149+
* @throws IOException if an error occurs
150+
*/
151+
public abstract FSDataOutputStream create(final Path f, WriteMode writeMode) throws IOException;
152+
138153
public abstract boolean delete(final Path f, final boolean recursive) throws IOException;
139154

140155
public abstract FileStatus[] listStatus(final Path f) throws IOException;

‎twister2/checkpointing/src/java/edu/iu/dsc/tws/checkpointing/client/CheckpointingClientImpl.java‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -45,12 +45,16 @@ public final class CheckpointingClientImpl implements MessageHandler, Checkpoint
4545

4646
private static final Logger LOG = Logger.getLogger(CheckpointingClientImpl.class.getName());
4747

48+
public static final String CONFIG_WAIT_TIME = "twister2.checkpointing.request.timeout";
49+
4850
private RRClient rrClient;
51+
private long waitTime;
4952
private Map<RequestID, Message> blockingResponse = new ConcurrentHashMap<>();
5053
private Map<RequestID, MessageHandler> asyncHandlers = new ConcurrentHashMap<>();
5154

52-
public CheckpointingClientImpl(RRClient rrClient) {
55+
public CheckpointingClientImpl(RRClient rrClient, long waitTime) {
5356
this.rrClient = rrClient;
57+
this.waitTime = waitTime;
5458
}
5559

5660
public void init() {
@@ -75,7 +79,7 @@ public Checkpoint.ComponentDiscoveryResponse sendDiscoveryMessage(
7579
.setFamily(family)
7680
.setIndex(index)
7781
.build(),
78-
10000
82+
this.waitTime
7983
);
8084
return (Checkpoint.ComponentDiscoveryResponse) this.blockingResponse.remove(requestID);
8185
}
@@ -93,7 +97,7 @@ public Checkpoint.FamilyInitializeResponse initFamily(int containerIndex,
9397
.setContainerIndex(containerIndex)
9498
.setContainers(containersCount)
9599
.build(),
96-
10000
100+
this.waitTime
97101
);
98102
return (Checkpoint.FamilyInitializeResponse) this.blockingResponse.remove(requestID);
99103
}

‎twister2/checkpointing/src/java/edu/iu/dsc/tws/checkpointing/master/CheckpointManager.java‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -178,7 +178,8 @@ private Long getFamilyVersion(String family) {
178178
}
179179

180180
private void handleFamilyInit(RequestID id, Checkpoint.FamilyInitialize message) {
181-
LOG.fine("Family init request received from " + message.getContainerIndex());
181+
LOG.fine("Family init request received from " + message.getContainerIndex()
182+
+ ". Family : " + message.getFamily());
182183

183184
FamilyInitHandler familyInitHandler = this.familyInitHandlers.get(message.getFamily());
184185

@@ -203,10 +204,13 @@ private void handleFamilyInit(RequestID id, Checkpoint.FamilyInitialize message)
203204
}
204205
}
205206

206-
boolean sentResponses = familyInitHandler.scheduleResponse(id);
207+
boolean sentResponses = familyInitHandler.scheduleResponse(message.getContainerIndex(), id);
207208
if (sentResponses) {
208209
LOG.info("Family " + message.getFamily() + " will start with version "
209210
+ familyInitHandler.getVersion());
211+
} else {
212+
LOG.fine("Scheduled family init response for family : " + message.getFamily()
213+
+ " for worker id " + message.getContainerIndex());
210214
}
211215
}
212216
}

‎twister2/checkpointing/src/java/edu/iu/dsc/tws/checkpointing/master/FamilyInitHandler.java‎

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,7 @@
1111
// limitations under the License.
1212
package edu.iu.dsc.tws.checkpointing.master;
1313

14-
import java.util.HashSet;
15-
import java.util.Set;
14+
import java.util.HashMap;
1615
import java.util.logging.Logger;
1716

1817
import 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;

‎twister2/config/src/yaml/conf/common/checkpoint.yaml‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,3 +13,6 @@ twister2.checkpointing.store.hdfs.dir: "/twister2/persistent/"
1313

1414
# Source triggering frequency
1515
twister2.checkpointing.source.frequency: 1000
16+
17+
# Checkpointing message request timeout
18+
twister2.checkpointing.request.timeout: 10000

‎twister2/config/src/yaml/conf/common/core.yaml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
# User name that will be used in JobID
33
###################################################################
44
# JobID is constructed as:
5-
# <username>-<jobName>-<timestamp>
5+
# [username]-[jobName]-[timestamp]
66
# if username is specified here, we use this value.
77
# Otherwise we get username from shell environment.
88
# if the username is longer than 9 characters, we use first 9 characters of it

‎twister2/config/src/yaml/conf/kubernetes/resource.yaml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -150,7 +150,7 @@ twister2.resource.class.uploader: "edu.iu.dsc.tws.rsched.uploaders.k8s.K8sUpload
150150

151151
# s3 bucket name to upload the job package
152152
# workers will download the job package from this location
153-
twister2.s3.bucket.name: "s3://<bucket-name>"
153+
twister2.s3.bucket.name: "s3://[bucket-name]"
154154

155155
# job package link will be available this much time
156156
# by default, it is 2 hours

‎twister2/data/src/main/java/edu/iu/dsc/tws/data/fs/local/LocalFileSystem.java‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -177,6 +177,12 @@ public FSDataOutputStream create(Path filePath) throws IOException {
177177
return new LocalDataOutputStream(file);
178178
}
179179

180+
@Override
181+
public FSDataOutputStream create(Path filePath, WriteMode writeMode) throws IOException {
182+
LOG.log(Level.WARNING, "This file system doesn't have the override mode implemented");
183+
return create(filePath);
184+
}
185+
180186
@Override
181187
public boolean delete(Path f, boolean recursive) throws IOException {
182188
final File file = pathToFile(f);

‎twister2/data/src/main/java/edu/iu/dsc/tws/data/hdfs/HadoopFileSystem.java‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import org.apache.hadoop.fs.RemoteIterator;
2424

2525
import edu.iu.dsc.tws.api.data.BlockLocation;
26+
import edu.iu.dsc.tws.api.data.FSDataOutputStream;
2627
import edu.iu.dsc.tws.api.data.FileStatus;
2728
import edu.iu.dsc.tws.api.data.FileSystem;
2829
import edu.iu.dsc.tws.api.data.Path;
@@ -149,6 +150,13 @@ public HadoopDataOutputStream create(final Path f) throws IOException {
149150
return new HadoopDataOutputStream(fsDataOutputStream);
150151
}
151152

153+
@Override
154+
public FSDataOutputStream create(Path f, WriteMode writeMode) throws IOException {
155+
final org.apache.hadoop.fs.FSDataOutputStream fsDataOutputStream =
156+
this.hadoopFileSystem.create(toHadoopPath(f), writeMode == WriteMode.OVERWRITE);
157+
return new HadoopDataOutputStream(fsDataOutputStream);
158+
}
159+
152160
public HadoopDataOutputStream append(Path path) throws IOException {
153161
final org.apache.hadoop.fs.FSDataOutputStream fsDataOutputStream =
154162
this.hadoopFileSystem.append(toHadoopPath(path));

0 commit comments

Comments
 (0)