You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@beam.apache.org by ke...@apache.org on 2016/08/22 19:09:22 UTC
[1/3] incubator-beam git commit: Remove unused constant in
ExecutorServiceParallelExecutor
Repository: incubator-beam
Updated Branches:
refs/heads/master 921d0b2e7 -> 064796dcc
Remove unused constant in ExecutorServiceParallelExecutor
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/b562072b
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/b562072b
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/b562072b
Branch: refs/heads/master
Commit: b562072b868bff14d7c33a64d05056c604710884
Parents: 921d0b2
Author: Thomas Groh <tg...@google.com>
Authored: Mon Aug 22 09:23:57 2016 -0700
Committer: Thomas Groh <tg...@google.com>
Committed: Mon Aug 22 09:23:57 2016 -0700
----------------------------------------------------------------------
.../beam/runners/direct/ExecutorServiceParallelExecutor.java | 2 --
1 file changed, 2 deletions(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/b562072b/runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java
----------------------------------------------------------------------
diff --git a/runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java b/runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java
index 8c6c6ed..35b6239 100644
--- a/runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java
+++ b/runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java
@@ -338,8 +338,6 @@ final class ExecutorServiceParallelExecutor implements PipelineExecutor {
}
private class MonitorRunnable implements Runnable {
- // arbitrary termination condition to ensure progress in the presence of pushback
- private final long maxTimeProcessingUpdatesNanos = TimeUnit.MILLISECONDS.toNanos(5L);
private final String runnableName = String.format("%s$%s-monitor",
evaluationContext.getPipelineOptions().getAppName(),
ExecutorServiceParallelExecutor.class.getSimpleName());
[2/3] incubator-beam git commit: Remove extra timer firings in
WatermarkManager
Posted by ke...@apache.org.
Remove extra timer firings in WatermarkManager
These timers should not be fired - the windows should be expired via the
GC timer, and any elements should be emitted if neccessary.
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/a65be9f9
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/a65be9f9
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/a65be9f9
Branch: refs/heads/master
Commit: a65be9f92c26ffbd00a72a644d619e636ba04b6e
Parents: b562072
Author: Thomas Groh <tg...@google.com>
Authored: Mon Aug 22 10:05:35 2016 -0700
Committer: Thomas Groh <tg...@google.com>
Committed: Mon Aug 22 10:05:35 2016 -0700
----------------------------------------------------------------------
.../apache/beam/runners/direct/WatermarkManager.java | 15 ++++-----------
1 file changed, 4 insertions(+), 11 deletions(-)
----------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/a65be9f9/runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java
----------------------------------------------------------------------
diff --git a/runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java b/runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java
index c8dfa8c..a44fa50 100644
--- a/runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java
+++ b/runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java
@@ -1139,17 +1139,10 @@ public class WatermarkManager {
inputWatermark.extractFiredEventTimeTimers();
Map<StructuralKey<?>, List<TimerData>> processingTimers;
Map<StructuralKey<?>, List<TimerData>> synchronizedTimers;
- if (inputWatermark.get().equals(BoundedWindow.TIMESTAMP_MAX_VALUE)) {
- processingTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
- TimeDomain.PROCESSING_TIME, BoundedWindow.TIMESTAMP_MAX_VALUE);
- synchronizedTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
- TimeDomain.PROCESSING_TIME, BoundedWindow.TIMESTAMP_MAX_VALUE);
- } else {
- processingTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
- TimeDomain.PROCESSING_TIME, clock.now());
- synchronizedTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
- TimeDomain.SYNCHRONIZED_PROCESSING_TIME, getSynchronizedProcessingInputTime());
- }
+ processingTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
+ TimeDomain.PROCESSING_TIME, clock.now());
+ synchronizedTimers = synchronizedProcessingInputWatermark.extractFiredDomainTimers(
+ TimeDomain.SYNCHRONIZED_PROCESSING_TIME, getSynchronizedProcessingInputTime());
Map<StructuralKey<?>, Map<TimeDomain, List<TimerData>>> groupedTimers = new HashMap<>();
groupFiredTimers(groupedTimers, eventTimeTimers, processingTimers, synchronizedTimers);
[3/3] incubator-beam git commit: This closes #861
Posted by ke...@apache.org.
This closes #861
Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/064796dc
Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/064796dc
Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/064796dc
Branch: refs/heads/master
Commit: 064796dcc6bffe5028411006c196129ee8766c44
Parents: 921d0b2 a65be9f
Author: Kenneth Knowles <kl...@google.com>
Authored: Mon Aug 22 11:15:40 2016 -0700
Committer: Kenneth Knowles <kl...@google.com>
Committed: Mon Aug 22 11:15:40 2016 -0700
----------------------------------------------------------------------
.../direct/ExecutorServiceParallelExecutor.java | 2 --
.../apache/beam/runners/direct/WatermarkManager.java | 15 ++++-----------
2 files changed, 4 insertions(+), 13 deletions(-)
----------------------------------------------------------------------