This is an automated email from the ASF dual-hosted git repository.
mattcasters pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new c482981821 Issue #8546 : Keep pipeline metrics updating after a Kafka
consumer restart (#8560)
c482981821 is described below
commit c482981821f62c738e0b6d7081e5baa08e4a51a8
Author: Matt Casters <[email protected]>
AuthorDate: Mon Sep 28 12:27:25 2026 +0200
Issue #8546 : Keep pipeline metrics updating after a Kafka consumer restart
(#8560)
* Issue #8546 : Keep pipeline metrics updating after a Kafka consumer
restart
A finished listener from the previous engine could stop the new run's
canvas and metrics timers, so the grid stayed at 0 until the next stop.
Kafka also logged each batch and then cleared the sub-pipeline counters.
* Issue #8546 : Address review on metrics timers and restart
Run again when the engine is not running, including a successful Beam
run and a failed preparation. Do GUI work outside the session lock, and
cancel a refresh timer once its engine is no longer on screen.
---
.../pages/pipeline/transforms/kafkaconsumer.adoc | 2 +
.../pipeline/SingleThreadedPipelineExecutor.java | 54 -------
.../kafka/consumer/KafkaConsumerInput.java | 17 ++-
.../kafka/consumer/KafkaConsumerInputTest.java | 9 ++
.../hopgui/file/pipeline/HopGuiPipelineGraph.java | 161 +++++++++++++-------
.../delegates/HopGuiPipelineGridDelegate.java | 56 ++++---
.../ui/hopgui/file/shared/ExecutionGuiSession.java | 163 +++++++++++++++++++++
.../hopgui/file/workflow/HopGuiWorkflowGraph.java | 104 ++++++++-----
.../file/shared/ExecutionGuiSessionTest.java | 122 +++++++++++++++
9 files changed, 512 insertions(+), 176 deletions(-)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
index 33baddc302..11bf38d132 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
@@ -27,6 +27,8 @@ The Kafka consumer transform pulls streaming data from Kafka.
It runs a sub-pipe
This sub-pipeline must start with an Injector transform.
+The Input metric on the Kafka consumer is the number of records consumed so
far in the run. Read and Written stay at 0 unless result fields are returned
from the sub-pipeline. Open the running sub-pipeline to see each transform's
own metrics. Those totals accumulate until the consumer stops.
+
You can define the number of messages to accept for processing, as well as the
specific data formats to stream activity data and system metrics.
You can set up this transform to collect monitored events, track user
consumption of data streams, and monitor alerts.
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
b/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
index e0b0d51259..5257f2017e 100644
---
a/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
+++
b/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
@@ -25,8 +25,6 @@ import org.apache.hop.core.IRowSet;
import org.apache.hop.core.Result;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.logging.ILogChannel;
-import org.apache.hop.i18n.BaseMessages;
-import org.apache.hop.pipeline.transform.BaseTransform;
import org.apache.hop.pipeline.transform.TransformMetaDataCombi;
import org.apache.hop.pipeline.transform.stream.IStream;
@@ -45,12 +43,9 @@ public class SingleThreadedPipelineExecutor {
private List<List<IStream>> transformInfoStreams;
private List<List<IRowSet>> transformInfoRowSets;
private ILogChannel log;
- private static final Class<?> PKG = SingleThreadedPipelineExecutor.class;
private static final String CONST_SEPARATOR =
"-------------------------------------------------------";
- @Getter @Setter private boolean clearingMetricsPerIteration = true;
-
public SingleThreadedPipelineExecutor(Pipeline pipeline) {
initializeObject(pipeline, false);
}
@@ -405,55 +400,6 @@ public class SingleThreadedPipelineExecutor {
return nrDone < transforms.size() && !pipeline.isStopped();
}
- public void buildExecutionSummary() {
-
- for (TransformMetaDataCombi combi : transforms) {
- // Summarize execution results
- long li = combi.transform.getLinesInput();
- long lo = combi.transform.getLinesOutput();
- long lr = combi.transform.getLinesRead();
- long lw = combi.transform.getLinesWritten();
- long lu = combi.transform.getLinesUpdated();
- long lj = combi.transform.getLinesRejected();
- long e = combi.transform.getErrors();
-
- ILogChannel tLog = combi.transform.getLogChannel();
-
- if (li > 0 || lo > 0 || lr > 0 || lw > 0 || lu > 0 || lj > 0 || e > 0) {
- tLog.logBasic(
- BaseMessages.getString(
- PKG,
- "SingleThreadedPipeline.Log.SummaryInfo",
- String.valueOf(li),
- String.valueOf(lo),
- String.valueOf(lr),
- String.valueOf(lw),
- String.valueOf(lu),
- String.valueOf(e + lj)));
- } else {
- tLog.logDetailed(
- BaseMessages.getString(
- PKG,
- "SingleThreadedPipeline.Log.SummaryInfo",
- String.valueOf(li),
- String.valueOf(lo),
- String.valueOf(lr),
- String.valueOf(lw),
- String.valueOf(lu),
- String.valueOf(e + lj)));
- }
- if (clearingMetricsPerIteration) {
- ((BaseTransform<?, ?>) combi.transform).setLinesInput(0);
- ((BaseTransform<?, ?>) combi.transform).setLinesOutput(0);
- ((BaseTransform<?, ?>) combi.transform).setLinesWritten(0);
- ((BaseTransform<?, ?>) combi.transform).setLinesRead(0);
- ((BaseTransform<?, ?>) combi.transform).setLinesSkipped(0);
- ((BaseTransform<?, ?>) combi.transform).setLinesUpdated(0);
- combi.transform.setLinesRejected(0);
- }
- }
- }
-
protected int getTotalRows(List<IRowSet> rowSets) {
int total = 0;
for (IRowSet rowSet : rowSets) {
diff --git
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
index 59cf04e850..3636b6a4f3 100644
---
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
+++
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
@@ -19,6 +19,7 @@ package org.apache.hop.pipeline.transforms.kafka.consumer;
import java.time.Duration;
import java.util.ArrayList;
+import java.util.Date;
import java.util.Properties;
import java.util.Set;
import java.util.stream.Collectors;
@@ -210,8 +211,6 @@ public class KafkaConsumerInput
// If the conditions for error handling are not met init
SingleThreadedExecutor normally
data.executor = new SingleThreadedPipelineExecutor(kafkaPipeline);
}
- data.executor.setClearingMetricsPerIteration(
- StringUtils.isEmpty(meta.getExecutionInformationLocation()));
// Initialize the sub-pipeline
//
@@ -351,6 +350,9 @@ public class KafkaConsumerInput
} else {
// Grab the records...
//
+ if (getFirstRowReadDate() == null) {
+ setFirstRowReadDate(new Date());
+ }
for (ConsumerRecord<Object, Object> record : records) {
Object[] outputRow = processMessageAsRow(record);
data.rowProducer.putRow(data.outputRowMeta, outputRow);
@@ -361,7 +363,7 @@ public class KafkaConsumerInput
}
data.lastRecordTime = System.currentTimeMillis();
if (isBasic()) {
- logBasic("Number of rows read: " +
data.rowProducer.getRowSet().size());
+ logBasic(batchLogMessage(records.count(), getLinesInput()));
}
// Pass them to the single threaded transformation and do an
iteration...
//
@@ -402,7 +404,6 @@ public class KafkaConsumerInput
// "removing" failing items from the kafka queue
//
data.consumer.commitAsync();
- data.executor.buildExecutionSummary();
if (errorHandlingConditionIsSatisfied()) {
data.incomingRowsBuffer.clear();
}
@@ -454,6 +455,14 @@ public class KafkaConsumerInput
return true;
}
+ /**
+ * One batch in the parent log. This is not the single-threaded executor's
"Finished processing"
+ * line, and it is not the injector buffer size.
+ */
+ static String batchLogMessage(int batchRecords, long totalInput) {
+ return "Kafka consumer batch of " + batchRecords + " record(s), cumulative
input " + totalInput;
+ }
+
/**
* True when a max consume duration is configured and the wall clock since
transform start has
* reached it. {@code maxConsumeDurationMs <= 0} means no limit.
diff --git
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
index d2fb5e186e..ac520e67f0 100644
---
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
+++
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
@@ -46,6 +46,15 @@ class KafkaConsumerInputTest {
assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(true, 2_000L, 0L, 0L,
0L));
}
+ @Test
+ void batchLogReportsBatchSizeAndCumulativeInput() {
+ String message = KafkaConsumerInput.batchLogMessage(12, 40L);
+ assertTrue(message.contains("12"));
+ assertTrue(message.contains("40"));
+ assertFalse(message.toLowerCase().contains("finished processing"));
+ assertFalse(message.toLowerCase().contains("rows read"));
+ }
+
@Test
void pollTimeoutIsCappedToRemainingConsumeDuration() {
// A deadline uses a short poll so empty topics re-check the clock instead
of blocking in
diff --git
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
index edf70fa058..634efa57c4 100644
---
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
+++
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
@@ -34,7 +34,6 @@ import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.Timer;
-import java.util.TimerTask;
import java.util.UUID;
import lombok.Getter;
import lombok.Setter;
@@ -200,6 +199,7 @@ import
org.apache.hop.ui.hopgui.file.pipeline.extension.HopGuiPipelineGraphExten
import
org.apache.hop.ui.hopgui.file.pipeline.extension.PipelineRenamedExtension;
import org.apache.hop.ui.hopgui.file.shared.CanvasToolTip;
import org.apache.hop.ui.hopgui.file.shared.DrillDownGuiPlugin;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
import org.apache.hop.ui.hopgui.file.shared.HopGuiAbstractGraph;
import org.apache.hop.ui.hopgui.file.shared.HopGuiGraphSnapshotUndo;
import org.apache.hop.ui.hopgui.file.shared.HopGuiTooltipExtension;
@@ -514,6 +514,9 @@ public class HopGuiPipelineGraph extends HopGuiAbstractGraph
Timer redrawTimer;
+ /** Ties redraw and metrics timers to the engine currently shown. */
+ private final ExecutionGuiSession executionGuiSession = new
ExecutionGuiSession();
+
@Setter private HopPipelineFileType<PipelineMeta> fileType;
private boolean doubleClick;
@@ -6269,6 +6272,8 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
if (handlePipelineMetaChanges(pipelineMeta)) {
// If the pipeline is not running, start the pipeline...
+ // Stopped counts as not running. Beam leaves a successful run
unfinished, and a failed
+ // preparation only stops the engine, so "not finished" would block the
next Run.
//
if (!isRunning()) {
try {
@@ -6303,12 +6308,12 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
//
pipelineMeta.clearCaches();
- pipeline =
+ setDisplayedPipeline(
PipelineEngineFactory.createPipelineEngine(
variables,
variables.resolve(pipelineRunConfigurationName),
hopGui.getMetadataProvider(),
- pipelineMeta);
+ pipelineMeta));
DrillDownGuiPlugin.bindToHopGui(pipeline, hopGui.getId());
// Set the variables from the execution configuration
@@ -6350,7 +6355,7 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
}
} catch (HopException e) {
- pipeline = null;
+ setDisplayedPipeline(null);
new ErrorDialog(
hopShell(),
BaseMessages.getString(PKG,
"PipelineLog.Dialog.ErrorOpeningPipeline.Title"),
@@ -6382,10 +6387,13 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
updateGui();
- // Update the GUI at the end of the pipeline
+ // Update the GUI at the end of the pipeline. Ignore the event when
this engine is no
+ // longer the one on screen (a restart overlapped its shutdown). The
work stays outside
+ // the session lock: it can take the graph lock or open a dialog.
//
- pipeline.addExecutionFinishedListener(e -> pipelineFinished());
- pipeline.addExecutionStoppedListener(e -> pipelineStopped());
+ final IPipelineEngine<PipelineMeta> engine = pipeline;
+ engine.addExecutionFinishedListener(this::runPipelineFinished);
+ engine.addExecutionStoppedListener(this::runPipelineStopped);
}
} else {
modalMessageDialog(
@@ -6503,7 +6511,8 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
// Create a new pipeline to execution
//
pipelineMeta.clearCaches();
- pipeline = new LocalPipelineEngine(pipelineMeta, variables,
hopGui.getLoggingObject());
+ setDisplayedPipeline(
+ new LocalPipelineEngine(pipelineMeta, variables,
hopGui.getLoggingObject()));
DrillDownGuiPlugin.bindToHopGui(pipeline, hopGui.getId());
pipeline.setPreview(true);
pipeline.setVariable(IPipelineEngine.PIPELINE_IN_PREVIEW_MODE, "Y");
@@ -6708,16 +6717,16 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
}
private synchronized void startThreads() {
+ final IPipelineEngine<PipelineMeta> engine = pipeline;
+ if (engine == null) {
+ return;
+ }
try {
// Add a listener to the pipeline.
// If the pipeline is done, we want to do the end processing, etc.
+ // A listener from an engine that is no longer on screen must not stop
the new run's timers.
//
- pipeline.addExecutionFinishedListener(
- p -> {
- checkPipelineEnded();
- checkErrorVisuals();
- stopRedrawTimer();
- });
+ engine.addExecutionFinishedListener(this::runPipelineMetricsFinished);
hopGui
.getDisplay()
@@ -6726,14 +6735,13 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
new Thread(
() -> {
try {
- pipeline.startThreads();
- pipeline.waitUntilFinished();
+ engine.startThreads();
+ engine.waitUntilFinished();
} catch (Exception e) {
- pipeline
+ engine
.getLogChannel()
.logError("Error starting transform
threads", e);
- checkErrorVisuals();
- stopRedrawTimer();
+ stopPipelineTimersIfCurrent(engine);
}
})
.start());
@@ -6742,13 +6750,8 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
updateGui();
} catch (Exception e) {
- if (pipeline != null) {
- pipeline.getLogChannel().logError("Error starting transform threads",
e);
- } else {
- log.logError("Error starting transform threads", e);
- }
- checkErrorVisuals();
- stopRedrawTimer();
+ engine.getLogChannel().logError("Error starting transform threads", e);
+ stopPipelineTimersIfCurrent(engine);
}
}
@@ -6763,8 +6766,10 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
return;
}
- // Set the pipeline instance
- this.pipeline = runningPipeline;
+ // Set the pipeline instance. Adopt it before timers and listeners so a
previous engine's
+ // finished event cannot own the refresh anymore.
+ //
+ setDisplayedPipeline(runningPipeline);
// Add all the execution result tabs (logging, metrics, etc.)
addAllTabs();
@@ -6794,14 +6799,9 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
if (isRunning) {
// Add listeners for when the pipeline finishes (only if still running)
- pipeline.addExecutionFinishedListener(
- p -> {
- checkPipelineEnded();
- checkErrorVisuals();
- stopRedrawTimer();
- });
+ pipeline.addExecutionFinishedListener(this::runPipelineMetricsFinished);
- pipeline.addExecutionStoppedListener(e -> pipelineStopped());
+ pipeline.addExecutionStoppedListener(this::runPipelineStopped);
// Start the redraw timer to continuously update the GUI
startRedrawTimer();
@@ -6817,28 +6817,79 @@ public class HopGuiPipelineGraph extends
HopGuiAbstractGraph
}
private void startRedrawTimer() {
-
- redrawTimer = new Timer("HopGuiPipelineGraph: redraw timer");
- TimerTask timtask =
- new TimerTask() {
- @Override
- public void run() {
- if (!hopDisplay().isDisposed()) {
- hopDisplay()
- .asyncExec(
- () -> {
- if (!HopGuiPipelineGraph.this.canvas.isDisposed()
- && perspective.isActive()
- && HopGuiPipelineGraph.this.isVisible()) {
- HopGuiPipelineGraph.this.canvas.redraw();
- HopGuiPipelineGraph.this.updateGui();
- }
- });
- }
+ ExecutionGuiSession.Snapshot snapshot = executionGuiSession.current();
+ if (!(snapshot.engine() instanceof IPipelineEngine<?> engine)) {
+ return;
+ }
+ executionGuiSession.scheduleWhileCurrent(
+ snapshot,
+ "HopGuiPipelineGraph: redraw timer",
+ ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH,
+ () -> !engine.isFinished(),
+ timer -> {
+ ExecutorUtil.cleanup(redrawTimer);
+ redrawTimer = timer;
+ },
+ () -> {
+ if (hopDisplay().isDisposed()) {
+ return;
}
- };
+ hopDisplay()
+ .asyncExec(
+ () -> {
+ if (!executionGuiSession.isCurrent(engine,
snapshot.generation())) {
+ return;
+ }
+ if (!HopGuiPipelineGraph.this.canvas.isDisposed()
+ && perspective.isActive()
+ && HopGuiPipelineGraph.this.isVisible()) {
+ HopGuiPipelineGraph.this.canvas.redraw();
+ HopGuiPipelineGraph.this.updateGui();
+ }
+ });
+ });
+ }
+
+ /**
+ * Publish {@code engine} as the pipeline on screen, together with the
session the timers check.
+ */
+ private void setDisplayedPipeline(IPipelineEngine<PipelineMeta> engine) {
+ executionGuiSession.adopt(engine, () -> this.pipeline = engine);
+ }
+
+ public ExecutionGuiSession getExecutionGuiSession() {
+ return executionGuiSession;
+ }
+
+ /** Extension-point finish hook. Runs outside the session lock because it
may open a dialog. */
+ private void runPipelineFinished(IPipelineEngine<PipelineMeta> finished) {
+ if (executionGuiSession.isCurrentEngine(finished)) {
+ pipelineFinished();
+ }
+ }
- redrawTimer.schedule(timtask, 0L,
ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH);
+ private void runPipelineStopped(IPipelineEngine<PipelineMeta> stopped) {
+ if (executionGuiSession.isCurrentEngine(stopped)) {
+ pipelineStopped();
+ }
+ }
+
+ /** Metrics refresh for a pipeline that just finished. The timer stop is the
only locked part. */
+ private void runPipelineMetricsFinished(IPipelineEngine<PipelineMeta>
finished) {
+ if (!executionGuiSession.isCurrentEngine(finished)) {
+ return;
+ }
+ checkPipelineEnded();
+ checkErrorVisuals();
+ executionGuiSession.stopIfCurrent(finished, this::stopRedrawTimer);
+ }
+
+ private void stopPipelineTimersIfCurrent(IPipelineEngine<PipelineMeta>
engine) {
+ if (!executionGuiSession.isCurrentEngine(engine)) {
+ return;
+ }
+ checkErrorVisuals();
+ executionGuiSession.stopIfCurrent(engine, this::stopRedrawTimer);
}
protected void stopRedrawTimer() {
diff --git
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
index bdfeab9325..f22b1aedaa 100644
---
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
+++
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
@@ -29,7 +29,6 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.Timer;
-import java.util.TimerTask;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
import lombok.Getter;
@@ -49,6 +48,7 @@ import org.apache.hop.pipeline.PipelineMeta;
import org.apache.hop.pipeline.engine.EngineMetrics;
import org.apache.hop.pipeline.engine.IEngineComponent;
import org.apache.hop.pipeline.engine.IEngineMetric;
+import org.apache.hop.pipeline.engine.IPipelineEngine;
import org.apache.hop.pipeline.transform.ITransform;
import org.apache.hop.pipeline.transform.TransformMeta;
import org.apache.hop.pipeline.transform.TransformStatus;
@@ -63,6 +63,7 @@ import org.apache.hop.ui.hopgui.ToolbarFacade;
import org.apache.hop.ui.hopgui.file.IHopFileTypeHandler;
import org.apache.hop.ui.hopgui.file.pipeline.HopGuiPipelineGraph;
import org.apache.hop.ui.hopgui.file.pipeline.PipelineMetricDisplayUtil;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
import org.eclipse.swt.SWT;
import org.eclipse.swt.custom.CLabel;
import org.eclipse.swt.custom.CTabItem;
@@ -327,31 +328,40 @@ public class HopGuiPipelineGridDelegate {
}
public void startRefreshMetricsTimer() {
- if (refreshMetricsTimer != null) {
+ // A new run replaces any timer left by the previous one. Returning early
when a timer already
+ // existed dropped the new run on the floor if the previous finish then
cancelled that timer.
+ //
+ ExecutionGuiSession session = pipelineGraph.getExecutionGuiSession();
+ ExecutionGuiSession.Snapshot snapshot = session.current();
+ if (!(snapshot.engine() instanceof IPipelineEngine<?> engine)) {
return;
}
-
- // Timer updates the view every UPDATE_TIME_VIEW interval
- refreshMetricsTimer = new Timer("HopGuiPipelineGraph: " +
pipelineGraph.getMeta().getName());
-
- TimerTask refreshMetricsTimerTask =
- new TimerTask() {
- @Override
- public void run() {
- if (!hopGui.getDisplay().isDisposed()) {
-
hopGui.getDisplay().asyncExec(HopGuiPipelineGridDelegate.this::refreshView);
- if (pipelineGraph.getPipeline() != null
- && (pipelineGraph.getPipeline().isFinished()
- || pipelineGraph.getPipeline().isStopped())
- && !pipelineGraph.getPipeline().isReadyToStart()) {
- ExecutorUtil.cleanup(refreshMetricsTimer, UPDATE_TIME_VIEW +
10);
- refreshMetricsTimer = null;
- }
- }
+ session.scheduleWhileCurrent(
+ snapshot,
+ "HopGuiPipelineGraph: " + pipelineGraph.getMeta().getName(),
+ UPDATE_TIME_VIEW,
+ null,
+ timer -> {
+ ExecutorUtil.cleanup(refreshMetricsTimer);
+ refreshMetricsTimer = timer;
+ },
+ () -> {
+ if (hopGui.getDisplay().isDisposed()) {
+ return;
}
- };
-
- refreshMetricsTimer.schedule(refreshMetricsTimerTask, 0L,
UPDATE_TIME_VIEW);
+ hopGui
+ .getDisplay()
+ .asyncExec(
+ () -> {
+ if (!session.isCurrent(engine, snapshot.generation())) {
+ return;
+ }
+ refreshView();
+ });
+ if ((engine.isFinished() || engine.isStopped()) &&
!engine.isReadyToStart()) {
+ session.stopIfCurrent(engine, this::stopRefreshMetricsTimer);
+ }
+ });
}
public void stopRefreshMetricsTimer() {
diff --git
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
new file mode 100644
index 0000000000..12b4a51f7e
--- /dev/null
+++
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
@@ -0,0 +1,163 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.ui.hopgui.file.shared;
+
+import java.util.Timer;
+import java.util.TimerTask;
+import java.util.function.BooleanSupplier;
+import java.util.function.Consumer;
+
+/**
+ * Binds GUI refresh timers to the execution currently on screen.
+ *
+ * <p>{@link #adopt(Object, Runnable)} and {@link #stopIfCurrent(Object,
Runnable)} share one lock,
+ * so a finished listener from the previous engine cannot stop the timers of
the engine that
+ * replaced it. {@code stop} and anything run from {@link
#scheduleWhileCurrent} must not take the
+ * graph lock or wait on the GUI thread. The GUI thread takes this lock while
adopting an engine,
+ * and a listener that waits for the GUI while holding it deadlocks.
+ */
+public final class ExecutionGuiSession {
+
+ /** One engine and the generation assigned to it, read under the session
lock. */
+ public record Snapshot(Object engine, int generation) {}
+
+ private final Object lock = new Object();
+ private Object engine;
+ private int generation;
+
+ /**
+ * Show {@code engine}. {@code bind} runs under the lock before the session
publishes the engine,
+ * so readers either see the previous pair or this one.
+ */
+ public Snapshot adopt(Object engine, Runnable bind) {
+ synchronized (lock) {
+ if (bind != null) {
+ bind.run();
+ }
+ this.engine = engine;
+ return new Snapshot(engine, ++generation);
+ }
+ }
+
+ public Snapshot adopt(Object engine) {
+ return adopt(engine, null);
+ }
+
+ /** Engine and generation as of one lock acquisition. */
+ public Snapshot current() {
+ synchronized (lock) {
+ return new Snapshot(engine, generation);
+ }
+ }
+
+ public boolean isCurrent(Object engine, int generation) {
+ synchronized (lock) {
+ return engine != null && this.engine == engine && this.generation ==
generation;
+ }
+ }
+
+ /** True when {@code engine} is the one on screen, ignoring generation. */
+ public boolean isCurrentEngine(Object engine) {
+ synchronized (lock) {
+ return engine != null && this.engine == engine;
+ }
+ }
+
+ /**
+ * Run {@code action} only while {@code engine} is still the one adopted at
{@code generation}.
+ * The lock is held for the whole action.
+ */
+ public boolean runIfCurrent(Object engine, int generation, Runnable action) {
+ synchronized (lock) {
+ if (engine == null || this.engine != engine || this.generation !=
generation) {
+ return false;
+ }
+ if (action != null) {
+ action.run();
+ }
+ return true;
+ }
+ }
+
+ /**
+ * Run {@code stop} only if {@code engine} is still the one on screen. The
lock is held across the
+ * check and {@code stop}. {@code stop} must not take the graph lock or wait
on the GUI thread.
+ */
+ public boolean stopIfCurrent(Object engine, Runnable stop) {
+ synchronized (lock) {
+ if (engine == null || this.engine != engine) {
+ return false;
+ }
+ if (stop != null) {
+ stop.run();
+ }
+ return true;
+ }
+ }
+
+ /**
+ * Replace the caller's timer with one that runs {@code onTick} while {@code
snapshot} is current.
+ * A tick that finds a newer engine cancels the timer, so a listener that no
longer owns the
+ * screen does not have to. {@code allow}, {@code replace} and {@code
onTick} run on the timer
+ * thread or under the session lock and must not wait on the GUI thread.
+ *
+ * @param allow checked under the lock before the timer is stored. When it
returns false, nothing
+ * is scheduled.
+ * @param replace stores the new timer and drops the previous one. Called
under the lock.
+ * @return false when nothing was scheduled
+ */
+ public boolean scheduleWhileCurrent(
+ Snapshot snapshot,
+ String threadName,
+ long periodMs,
+ BooleanSupplier allow,
+ Consumer<Timer> replace,
+ Runnable onTick) {
+ if (snapshot == null || snapshot.engine() == null || replace == null ||
onTick == null) {
+ return false;
+ }
+ Timer timer = new Timer(threadName);
+ TimerTask task =
+ new TimerTask() {
+ @Override
+ public void run() {
+ if (!isCurrent(snapshot.engine(), snapshot.generation())) {
+ timer.cancel();
+ return;
+ }
+ onTick.run();
+ }
+ };
+ boolean[] started = {false};
+ runIfCurrent(
+ snapshot.engine(),
+ snapshot.generation(),
+ () -> {
+ if (allow != null && !allow.getAsBoolean()) {
+ return;
+ }
+ replace.accept(timer);
+ timer.schedule(task, 0L, periodMs);
+ started[0] = true;
+ });
+ if (!started[0]) {
+ timer.cancel();
+ }
+ return started[0];
+ }
+}
diff --git
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
index f3c38ac91a..a45c2cca1c 100644
---
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
+++
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
@@ -29,7 +29,6 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.Timer;
-import java.util.TimerTask;
import java.util.UUID;
import lombok.Getter;
import lombok.Setter;
@@ -139,6 +138,7 @@ import
org.apache.hop.ui.hopgui.file.delegates.HopGuiNoteLinkSupport;
import org.apache.hop.ui.hopgui.file.delegates.HopGuiNotePadDelegate;
import org.apache.hop.ui.hopgui.file.shared.CanvasToolTip;
import org.apache.hop.ui.hopgui.file.shared.DrillDownGuiPlugin;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
import org.apache.hop.ui.hopgui.file.shared.HopGuiAbstractGraph;
import org.apache.hop.ui.hopgui.file.shared.HopGuiGraphSnapshotUndo;
import org.apache.hop.ui.hopgui.file.shared.HopGuiTooltipExtension;
@@ -445,6 +445,9 @@ public class HopGuiWorkflowGraph extends HopGuiAbstractGraph
private Timer redrawTimer;
+ /** Ties the canvas redraw timer to the workflow currently shown. */
+ private final ExecutionGuiSession executionGuiSession = new
ExecutionGuiSession();
+
public HopGuiWorkflowGraph(
Composite parent,
final HopGui hopGui,
@@ -4218,7 +4221,7 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
}
public synchronized void setWorkflow(IWorkflowEngine<WorkflowMeta> workflow)
{
- this.workflow = workflow;
+ executionGuiSession.adopt(workflow, () -> this.workflow = workflow);
}
public void paintControl(PaintEvent e) {
@@ -5423,13 +5426,13 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
ExtensionPointHandler.callExtensionPoint(
log, variables,
HopExtensionPoint.HopGuiWorkflowMetaExecutionStart.id, workflowMeta);
- workflow =
+ setWorkflow(
WorkflowEngineFactory.createWorkflowEngine(
variables,
variables.resolve(executionConfiguration.getRunConfiguration()),
hopGui.getMetadataProvider(),
runWorkflowMeta,
- hopGuiLoggingObject);
+ hopGuiLoggingObject));
workflow.setLogLevel(executionConfiguration.getLogLevel());
workflow.setGatheringMetrics(executionConfiguration.isGatheringMetrics());
@@ -5485,6 +5488,13 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
}
log.logBasic(BaseMessages.getString(PKG,
"WorkflowLog.Log.StartingWorkflow"));
+
+ // Listeners before the thread and the timer. A workflow that
finishes in between would
+ // otherwise leave the redraw timer running.
+ //
+ workflow.addExecutionFinishedListener(this::onWorkflowFinished);
+ workflow.addExecutionStoppedListener(this::onWorkflowStopped);
+
workflowThread = new Thread(() -> workflow.startExecution());
workflowThread.start();
@@ -5494,11 +5504,6 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
startRedrawTimer();
updateGui();
-
- // Attach a listener to notify us that the workflow has finished.
- //
- workflow.addExecutionFinishedListener(e ->
HopGuiWorkflowGraph.this.workflowFinished());
- workflow.addExecutionStoppedListener(e ->
HopGuiWorkflowGraph.this.workflowStopped());
// Show the execution results views
//
addAllTabs();
@@ -5508,7 +5513,7 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
BaseMessages.getString(PKG,
"WorkflowLog.Dialog.CanNotOpenWorkflow.Title"),
BaseMessages.getString(PKG,
"WorkflowLog.Dialog.CanNotOpenWorkflow.Message"),
e);
- workflow = null;
+ setWorkflow(null);
}
} else {
MessageBox m = new MessageBox(hopShell(), SWT.OK | SWT.ICON_WARNING);
@@ -5551,27 +5556,36 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
/** This gets called at the very end, when everything is done. */
protected void workflowFinished() {
- // Do a final check to see if it all ended...
- //
- if (workflow != null && workflow.isInitialized() && workflow.isFinished())
{
+ onWorkflowFinished(workflow);
+ }
+
+ private void onWorkflowFinished(IWorkflowEngine<WorkflowMeta> finished) {
+ if (!executionGuiSession.isCurrentEngine(finished)) {
+ return;
+ }
+ if (finished.isInitialized() && finished.isFinished()) {
log.logBasic(
BaseMessages.getString(PKG, "WorkflowLog.Log.WorkflowHasEnded",
workflowMeta.getName()));
}
-
- stopRedrawTimer();
-
updateGui();
+ executionGuiSession.stopIfCurrent(finished, this::stopRedrawTimer);
}
protected void workflowStopped() {
- if (workflow != null && workflow.isInitialized() && workflow.isStopped()) {
+ onWorkflowStopped(workflow);
+ }
+
+ private void onWorkflowStopped(IWorkflowEngine<WorkflowMeta> stopped) {
+ if (!executionGuiSession.isCurrentEngine(stopped)) {
+ return;
+ }
+ if (stopped.isInitialized() && stopped.isStopped()) {
log.logBasic(
BaseMessages.getString(
PKG, "WorkflowLog.Log.ProcessingOfWorkflowStopped",
workflowMeta.getName()));
}
-
- stopRedrawTimer();
updateGui();
+ executionGuiSession.stopIfCurrent(stopped, this::stopRedrawTimer);
}
@Override
@@ -6068,9 +6082,9 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
if (isRunning) {
// Add listeners for when the workflow finishes (only if still running)
- workflow.addExecutionFinishedListener(e ->
HopGuiWorkflowGraph.this.workflowFinished());
+ workflow.addExecutionFinishedListener(this::onWorkflowFinished);
- workflow.addExecutionStoppedListener(e ->
HopGuiWorkflowGraph.this.workflowStopped());
+ workflow.addExecutionStoppedListener(this::onWorkflowStopped);
// Start the redraw timer to continuously update the GUI
startRedrawTimer();
@@ -6085,26 +6099,36 @@ public class HopGuiWorkflowGraph extends
HopGuiAbstractGraph
}
private void startRedrawTimer() {
- redrawTimer = new Timer("WorkflowGraph auto refresh: " +
workflow.getWorkflowName());
- TimerTask timerTask =
- new TimerTask() {
- @Override
- public void run() {
- if (!hopDisplay().isDisposed()) {
- hopDisplay()
- .asyncExec(
- () -> {
- if (!HopGuiWorkflowGraph.this.canvas.isDisposed()
- && perspective.isActive()
- && HopGuiWorkflowGraph.this.isVisible()) {
- updateGui();
- }
- });
- }
+ ExecutionGuiSession.Snapshot snapshot = executionGuiSession.current();
+ if (!(snapshot.engine() instanceof IWorkflowEngine<?> engine)) {
+ return;
+ }
+ executionGuiSession.scheduleWhileCurrent(
+ snapshot,
+ "WorkflowGraph auto refresh: " + engine.getWorkflowName(),
+ ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH,
+ () -> !engine.isFinished(),
+ timer -> {
+ ExecutorUtil.cleanup(redrawTimer);
+ redrawTimer = timer;
+ },
+ () -> {
+ if (hopDisplay().isDisposed()) {
+ return;
}
- };
-
- redrawTimer.schedule(timerTask, 0L,
ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH);
+ hopDisplay()
+ .asyncExec(
+ () -> {
+ if (!executionGuiSession.isCurrent(engine,
snapshot.generation())) {
+ return;
+ }
+ if (!HopGuiWorkflowGraph.this.canvas.isDisposed()
+ && perspective.isActive()
+ && HopGuiWorkflowGraph.this.isVisible()) {
+ updateGui();
+ }
+ });
+ });
}
protected void stopRedrawTimer() {
diff --git
a/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
b/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
new file mode 100644
index 0000000000..9417ed76c2
--- /dev/null
+++
b/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
@@ -0,0 +1,122 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.ui.hopgui.file.shared;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.junit.jupiter.api.Test;
+
+class ExecutionGuiSessionTest {
+
+ @Test
+ void stopIfCurrentRunsForTheEngineOnScreen() {
+ ExecutionGuiSession session = new ExecutionGuiSession();
+ Object engine = new Object();
+ ExecutionGuiSession.Snapshot snapshot = session.adopt(engine);
+ AtomicBoolean stopped = new AtomicBoolean(false);
+
+ assertTrue(session.stopIfCurrent(engine, () -> stopped.set(true)));
+ assertTrue(stopped.get());
+ assertTrue(session.isCurrent(snapshot.engine(), snapshot.generation()));
+ assertTrue(session.isCurrentEngine(engine));
+ }
+
+ @Test
+ void stopIfCurrentIgnoresTheEngineThatWasReplaced() {
+ ExecutionGuiSession session = new ExecutionGuiSession();
+ Object first = new Object();
+ Object second = new Object();
+ session.adopt(first);
+ int secondGeneration = session.adopt(second).generation();
+ AtomicBoolean stopped = new AtomicBoolean(false);
+
+ assertFalse(session.stopIfCurrent(first, () -> stopped.set(true)));
+ assertFalse(stopped.get());
+ assertTrue(session.isCurrent(second, secondGeneration));
+ assertFalse(session.runIfCurrent(first, secondGeneration, () ->
stopped.set(true)));
+ assertFalse(stopped.get());
+ }
+
+ @Test
+ void adoptFromAnotherThreadWaitsUntilTheStopperReleasesTheLock() throws
Exception {
+ ExecutionGuiSession session = new ExecutionGuiSession();
+ Object first = new Object();
+ Object second = new Object();
+ Object third = new Object();
+ session.adopt(first);
+
+ CountDownLatch holding = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ AtomicInteger secondGeneration = new AtomicInteger();
+ AtomicBoolean failed = new AtomicBoolean(false);
+ Thread stopper =
+ new Thread(
+ () -> {
+ try {
+ session.stopIfCurrent(
+ first,
+ () -> {
+ secondGeneration.set(session.adopt(second).generation());
+ holding.countDown();
+ try {
+ if (!release.await(5, TimeUnit.SECONDS)) {
+ throw new AssertionError("stopper was not released");
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new AssertionError(e);
+ }
+ });
+ } catch (AssertionError e) {
+ failed.set(true);
+ }
+ });
+ stopper.start();
+ assertTrue(holding.await(5, TimeUnit.SECONDS));
+
+ CountDownLatch reachedAdopt = new CountDownLatch(1);
+ CountDownLatch adopted = new CountDownLatch(1);
+ AtomicInteger thirdGeneration = new AtomicInteger();
+ Thread adopter =
+ new Thread(
+ () -> {
+ reachedAdopt.countDown();
+ thirdGeneration.set(session.adopt(third).generation());
+ adopted.countDown();
+ });
+ adopter.start();
+ assertTrue(reachedAdopt.await(5, TimeUnit.SECONDS));
+ // Do not call into the session here: the stopper still holds its lock and
is waiting for
+ // release. The adopted latch stays open until that lock is released.
+ assertFalse(adopted.await(300, TimeUnit.MILLISECONDS));
+
+ release.countDown();
+ assertTrue(adopted.await(5, TimeUnit.SECONDS));
+ assertTrue(thirdGeneration.get() > secondGeneration.get());
+ assertTrue(session.isCurrent(third, thirdGeneration.get()));
+ assertFalse(session.isCurrent(second, secondGeneration.get()));
+ stopper.join(5_000);
+ assertFalse(stopper.isAlive());
+ assertFalse(failed.get());
+ }
+}