forked from apache/flink
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge branch 'version011' into version02
- Loading branch information
Showing
25 changed files
with
712 additions
and
477 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -33,6 +33,7 @@ | |
|
||
import eu.stratosphere.nephele.util.CommonTestUtils; | ||
|
||
|
||
/** | ||
* @author Mathias Peters <[email protected]> | ||
* TODO: {@link StringRecord} has a lot of public methods that need to be tested. | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
50 changes: 50 additions & 0 deletions
50
...udmanager/src/main/java/eu/stratosphere/nephele/instance/cloud/CloudInstanceNotifier.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,50 @@ | ||
package eu.stratosphere.nephele.instance.cloud; | ||
|
||
import eu.stratosphere.nephele.instance.AbstractInstance; | ||
import eu.stratosphere.nephele.instance.AllocatedResource; | ||
import eu.stratosphere.nephele.instance.InstanceListener; | ||
import eu.stratosphere.nephele.jobgraph.JobID; | ||
|
||
/** | ||
* This class is an auxiliary class to send the notification | ||
* about the availability of an {@link AbstractInstance} to the given {@link InstanceListener} object. The notification | ||
* must be sent from | ||
* a separate thread, otherwise the atomic operation of requesting an instance | ||
* for a vertex and switching to the state ASSINING could not be guaranteed. | ||
* This class is thread-safe. | ||
* | ||
* @author warneke | ||
*/ | ||
public class CloudInstanceNotifier extends Thread { | ||
|
||
/** | ||
* The {@link InstanceListener} object to send the notification to. | ||
*/ | ||
private final InstanceListener instanceListener; | ||
|
||
private final CloudInstance instance; | ||
|
||
private final JobID id; | ||
|
||
/** | ||
* Constructs a new instance notifier object. | ||
* | ||
* @param instanceListener | ||
* the listener to send the notification to | ||
* @param allocatedSlice | ||
* the slice with has been allocated for the job | ||
*/ | ||
public CloudInstanceNotifier(InstanceListener instanceListener, JobID id, CloudInstance instance) { | ||
this.instanceListener = instanceListener; | ||
this.instance = instance; | ||
this.id = id; | ||
} | ||
|
||
/** | ||
* {@inheritDoc} | ||
*/ | ||
@Override | ||
public void run() { | ||
this.instanceListener.resourceAllocated(id, instance.asAllocatedResource()); | ||
} | ||
} |
Oops, something went wrong.