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 8ecd3802d6 Fix memory leak in caching execution information location,
fixes #8623 (#8625)
8ecd3802d6 is described below
commit 8ecd3802d6df1590d83ff273ef8246170d48ff8b
Author: Matt Casters <[email protected]>
AuthorDate: Mon Sep 28 10:43:34 2026 +0200
Fix memory leak in caching execution information location, fixes #8623
(#8625)
* Fix memory leak in caching execution information location, fixes #8623
- Bound in-memory cache to 50 LRU entries by default with LinkedHashMap
- Reduce default maxCacheAge to 10 minutes (600000 ms)
- Add maxCacheSize and update maxCacheAge GUI widgets and i18n
- Enforce LRU eviction in BaseCachingExecutionInfoLocation and clear cache
on close()
- Fix CacheEntry.isTooOld() to eliminate immortal unwritten entries
- Fix query and PreparedStatement collision/leak in
CachingDatabaseExecutionInfoLocation.retrieveIds()
- Reuse Jackson ObjectMapper with HopJson.newMapper()
- Protect Database.insertRow() with try-finally closeInsert()
- Harden Pipeline and LocalPipelineEngine lifecycle so timers and locations
are always closed on completion, abort, or error
* Issue #8623 : Cap cached execution logs and keep a failed cache write
The live pipeline entry copied the whole log buffer on every tick. Store
only the new lines and keep the newest 2 million characters. Leave an
entry in memory when its save fails, and do not close the location on a
safe stop. A missing max cache age stays at one day; new locations use
10 minutes.
* Issue #8623 : Do not track single-threaded mappings in execution info
A single-threaded parent drives the mapping again on every iteration, and
the child does not finish between batches. A location on that child kept a
second caching session open for the whole parent run. Drop the location
from a copy of the run configuration, and stop the child on dispose.
* Issue #8623 : Do not rewrite the cached execution document on every save
* Issue #8623 : Keep execution logs incremental and finish before closing
A full log snapshot is no longer appended on top of the lines already
stored.
Engines that asked for every line now pass the previous line number. A
normal
pipeline waits for its finished state before the location is closed. A
single-threaded child still closes on stop, because it never reaches
pipelineCompleted(). The database location drops its connection when the
last
flush fails, and it no longer alters an existing table to add state_json.
* Issue #8623 : Keep the execution document current on a table without the
state column
A light update writes the execution state to state_json and leaves the
inserted document alone. A table created before that column has nowhere
else to keep the state, so the status columns moved while the document
kept the state and the log of the first save. Such a table now rewrites
the whole document on every save, as before, and keeps the metadata and
the pipeline XML in the entry so the rewrite does not drop them. The
statement that adds the column is logged when the column is missing.
---------
Co-authored-by: Bart Maertens <[email protected]>
---
.../org/apache/hop/core/database/Database.java | 9 +-
.../hop/execution/ExecutionStateBuilder.java | 20 +-
.../caching/BaseCachingExecutionInfoLocation.java | 223 +++++++++---
.../apache/hop/execution/caching/CacheEntry.java | 36 +-
.../sampler/ExecutionDataSamplerStoreBase.java | 7 +
.../sampler/IExecutionDataSamplerStore.java | 3 +
.../BasicDataProfilingDataSamplerStore.java | 7 +
.../java/org/apache/hop/pipeline/Pipeline.java | 56 +--
.../loadbalance/LoadBalancingPipelineEngine.java | 16 +-
.../engines/local/LocalPipelineEngine.java | 235 ++++++++-----
.../loadbalance/LoadBalancingWorkflowEngine.java | 16 +-
.../engines/local/LocalWorkflowEngine.java | 14 +-
.../caching/messages/messages_en_US.properties | 4 +-
.../caching/messages/messages_pt_BR.properties | 4 +-
.../BaseCachingExecutionInfoLocationTest.java | 231 +++++++++++++
.../local/LocalPipelineEngineExecutionIdTest.java | 172 ++++++++++
.../hop/beam/engines/BeamPipelineEngine.java | 21 +-
.../dataflow/BeamDataFlowPipelineEngine.java | 3 +-
.../hop/spark/engines/SparkPipelineEngine.java | 25 +-
.../CachingDatabaseExecutionInfoLocation.java | 376 ++++++++++++++++-----
.../CachingDatabaseExecutionInfoLocationTest.java | 317 +++++++++++++++++
.../kafka/consumer/KafkaConsumerInput.java | 5 +-
.../pipeline/transforms/mapping/SimpleMapping.java | 32 ++
.../mapping/messages/messages_en_US.properties | 1 +
.../mapping/SimpleMappingExecutionInfoTest.java | 60 ++++
.../ui/execution/ExecutionInfoLocationEditor.java | 7 +
26 files changed, 1648 insertions(+), 252 deletions(-)
diff --git a/core/src/main/java/org/apache/hop/core/database/Database.java
b/core/src/main/java/org/apache/hop/core/database/Database.java
index 63d5fda9b2..d426640056 100644
--- a/core/src/main/java/org/apache/hop/core/database/Database.java
+++ b/core/src/main/java/org/apache/hop/core/database/Database.java
@@ -1214,9 +1214,12 @@ public class Database implements IVariables,
ILoggingObject, AutoCloseable {
public void insertRow(String schemaName, String tableName, IRowMeta fields,
Object[] data)
throws HopDatabaseException {
prepareInsert(fields, schemaName, tableName);
- setValuesInsert(fields, data);
- insertRow();
- closeInsert();
+ try {
+ setValuesInsert(fields, data);
+ insertRow();
+ } finally {
+ closeInsert();
+ }
}
public String getInsertStatement(String tableName, IRowMeta fields) {
diff --git
a/engine/src/main/java/org/apache/hop/execution/ExecutionStateBuilder.java
b/engine/src/main/java/org/apache/hop/execution/ExecutionStateBuilder.java
index 0b803d316a..07e5552b17 100644
--- a/engine/src/main/java/org/apache/hop/execution/ExecutionStateBuilder.java
+++ b/engine/src/main/java/org/apache/hop/execution/ExecutionStateBuilder.java
@@ -71,9 +71,18 @@ public final class ExecutionStateBuilder {
return new ExecutionStateBuilder();
}
+ /**
+ * A null or negative request is the whole buffer. Callers pass {@code -1}
for that. A caching
+ * location appends when the state carries a line number, so a full snapshot
has to leave it unset
+ * or the next tick stores the same lines again.
+ */
+ private static boolean isDeltaRequest(Integer lastLogLineNr) {
+ return lastLogLineNr != null && lastLogLineNr >= 0;
+ }
+
private static String getLoggingText(String logChannelId, Integer
lastLogLineNr) {
StringBuffer loggingTextBuffer;
- if (lastLogLineNr != null) {
+ if (isDeltaRequest(lastLogLineNr)) {
loggingTextBuffer = HopLogStore.getAppender().getBuffer(logChannelId,
false, lastLogLineNr);
} else {
loggingTextBuffer = HopLogStore.getAppender().getBuffer(logChannelId,
false);
@@ -81,6 +90,11 @@ public final class ExecutionStateBuilder {
return loggingTextBuffer.toString();
}
+ /** Line number for the next delta. A full snapshot does not advance a
cursor. */
+ private static Integer loggingCursor(Integer requestedLineNr, int
lastNrInLogStore) {
+ return isDeltaRequest(requestedLineNr) ? lastNrInLogStore : null;
+ }
+
public static ExecutionStateBuilder fromExecutor(
IPipelineEngine<PipelineMeta> pipeline, Integer lastLogLineNr) {
String parentLogChannelId =
@@ -97,7 +111,7 @@ public final class ExecutionStateBuilder {
.withId(pipeline.getLogChannelId())
.withName(pipeline.getPipelineMeta().getName())
.withLoggingText(getLoggingText(pipeline.getLogChannelId(),
lastLogLineNr))
- .withLastLogLineNr(lastNrInLogStore)
+ .withLastLogLineNr(loggingCursor(lastLogLineNr, lastNrInLogStore))
.withFailed(pipeline.getErrors() > 0)
.withStatusDescription(pipeline.getStatusDescription())
.withChildIds(
@@ -261,7 +275,7 @@ public final class ExecutionStateBuilder {
.withId(workflow.getLogChannelId())
.withName(workflow.getWorkflowMeta().getName())
.withLoggingText(getLoggingText(workflow.getLogChannelId(),
lastLogLineNr))
- .withLastLogLineNr(lastNrInLogStore)
+ .withLastLogLineNr(loggingCursor(lastLogLineNr, lastNrInLogStore))
.withFailed(result != null && !result.isResult())
.withStatusDescription(workflow.getStatusDescription())
.withChildIds(
diff --git
a/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
b/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
index dc69ef19e7..a4f0de413b 100644
---
a/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
+++
b/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
@@ -23,8 +23,8 @@ import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.Date;
-import java.util.HashMap;
import java.util.HashSet;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -67,6 +67,31 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
@HopMetadataProperty
protected String persistenceDelay = "5000";
+ @GuiWidgetElement(
+ id = "maxCacheSize",
+ order = "905",
+ parentId = ExecutionInfoLocation.GUI_PLUGIN_ELEMENT_PARENT_ID,
+ type = GuiElementType.TEXT,
+ toolTip = "i18n::CachingFileExecutionInfoLocation.MaxCacheSize.Tooltip",
+ label = "i18n::CachingFileExecutionInfoLocation.MaxCacheSize.Label")
+ @HopMetadataProperty
+ protected String maxCacheSize = "50";
+
+ /**
+ * Locations saved before {@code maxCacheAge} existed omit the property. The
metadata loader
+ * leaves this initializer in place, so those locations keep the original
one-day age.
+ */
+ public static final String LEGACY_MAX_CACHE_AGE = "86400000";
+
+ /** Age applied when a caching location is created in the GUI, not when an
old file is loaded. */
+ public static final String NEW_LOCATION_MAX_CACHE_AGE = "600000";
+
+ /**
+ * Logging text kept on one cache entry. Matches the execution viewer's
default display limit and
+ * keeps the newest lines, so a long run cannot retain the whole log buffer
here as well.
+ */
+ public static final int MAX_CACHED_LOGGING_TEXT_CHARS = 2_000_000;
+
@GuiWidgetElement(
id = "maxCacheAge",
order = "910",
@@ -75,7 +100,7 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
toolTip = "i18n::CachingFileExecutionInfoLocation.MaxCacheAge.Tooltip",
label = "i18n::CachingFileExecutionInfoLocation.MaxCacheAge.Label")
@HopMetadataProperty
- protected String maxCacheAge = "86400000";
+ protected String maxCacheAge = LEGACY_MAX_CACHE_AGE;
protected IVariables variables;
protected IHopMetadataProvider metadataProvider;
@@ -105,21 +130,24 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
protected int delay;
protected int maxAge;
+ protected int maxSize;
protected BaseCachingExecutionInfoLocation() {
- cache = new HashMap<>();
+ cache = new LinkedHashMap<>(16, 0.75f, true);
this.cacheTimer = null;
this.locked = new AtomicBoolean(false);
}
protected BaseCachingExecutionInfoLocation(BaseCachingExecutionInfoLocation
location) {
this();
+ this.maxCacheSize = location.maxCacheSize;
this.maxCacheAge = location.maxCacheAge;
this.persistenceDelay = location.persistenceDelay;
this.variables = location.variables;
this.metadataProvider = location.metadataProvider;
this.delay = location.delay;
this.maxAge = location.maxAge;
+ this.maxSize = location.maxSize;
}
public abstract BaseCachingExecutionInfoLocation clone();
@@ -144,9 +172,16 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
//
delay = Const.toInt(variables.resolve(persistenceDelay), 60000);
- // The default maximum cache age is 1 day
+ // The default maximum cache size is 50
+ //
+ maxSize = Const.toInt(variables.resolve(maxCacheSize), 50);
+ if (maxSize <= 0) {
+ maxSize = 50;
+ }
+
+ // A missing age is the pre-existing 1 day. New GUI locations save 10
minutes explicitly.
//
- maxAge = Const.toInt(variables.resolve(maxCacheAge), 86400000);
+ maxAge = Const.toInt(variables.resolve(maxCacheAge),
Integer.parseInt(LEGACY_MAX_CACHE_AGE));
// Let's start a timer to manage the cache every second or so.
// Cancel any previous timer first: a second initialize() used to leave
the old one running
@@ -192,7 +227,12 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
//
for (CacheEntry cacheEntry : cache.values()) {
if (cacheEntry.needsWriting(delay)) {
- persistCacheEntry(cacheEntry);
+ try {
+ persistCacheEntry(cacheEntry);
+ } catch (Exception e) {
+ LogChannel.GENERAL.logError(
+ "Error persisting cache entry for " + cacheEntry.getId(), e);
+ }
}
}
@@ -207,6 +247,9 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
// Remove these entries.
//
tooOld.forEach(id -> cache.remove(id));
+
+ enforceMaxCacheSize();
+
loggedManageCacheError = false;
lastManageCacheError = null;
} catch (Exception e) {
@@ -222,6 +265,29 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
}
+ protected synchronized void enforceMaxCacheSize() {
+ int max = maxSize > 0 ? maxSize : 50;
+ if (cache.size() <= max) {
+ return;
+ }
+ var iterator = cache.entrySet().iterator();
+ while (iterator.hasNext() && cache.size() > max) {
+ Map.Entry<String, CacheEntry> entry = iterator.next();
+ CacheEntry cacheEntry = entry.getValue();
+ if (cacheEntry != null && cacheEntry.isDirty()) {
+ try {
+ persistCacheEntry(cacheEntry);
+ } catch (Exception e) {
+ // A failed write must stay in memory. Dropping it here loses the
only copy.
+ LogChannel.GENERAL.logError(
+ "Error persisting cache entry during eviction: " +
cacheEntry.getId(), e);
+ continue;
+ }
+ }
+ iterator.remove();
+ }
+ }
+
@Override
public synchronized void close() throws HopException {
// Stop the timer before the final flush. Tasks that are already inside
manageCache hold this
@@ -230,22 +296,31 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
Timer timer = cacheTimer;
cacheTimer = null;
ExecutorUtil.cleanup(timer);
- try {
- for (CacheEntry cacheEntry : cache.values()) {
- if (cacheEntry.isDirty()) {
+ HopException failure = null;
+ for (Map.Entry<String, CacheEntry> mapEntry : new
ArrayList<>(cache.entrySet())) {
+ CacheEntry cacheEntry = mapEntry.getValue();
+ if (cacheEntry != null && cacheEntry.isDirty()) {
+ try {
persistCacheEntry(cacheEntry);
+ } catch (Exception e) {
+ // Leave this entry. A later close() can retry it. Clearing it drops
unsaved state.
+ if (failure == null) {
+ failure =
+ new HopException("Error persisting caching execution
information location", e);
+ }
+ continue;
}
}
- } catch (Exception e) {
- throw new HopException("Error persisting caching execution information
location", e);
+ cache.remove(mapEntry.getKey());
+ }
+ if (failure != null) {
+ throw failure;
}
}
@Override
- public void clearCaches() {
- synchronized (locked) {
- cache.clear();
- }
+ public synchronized void clearCaches() {
+ cache.clear();
}
@Override
@@ -357,6 +432,7 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
entry.setLastWritten(null);
cache.put(execution.getId(), entry);
+ enforceMaxCacheSize();
}
protected synchronized void addChildExecutionToCache(Execution execution)
throws HopException {
@@ -385,6 +461,35 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
}
+ /**
+ * Pipeline and workflow updates carry only the lines written since {@code
lastLogLineNr}. Append
+ * that delta and keep the newest characters. A full snapshot ({@code
lastLogLineNr == null},
+ * including a caller that asked for every line with {@code -1}) replaces
the stored text. It is
+ * capped the same way so the cache does not keep a second copy of the
central log buffer.
+ */
+ static void appendLoggingDelta(ExecutionState previous, ExecutionState
update) {
+ if (update == null) {
+ return;
+ }
+ String delta = capLoggingText(update.getLoggingText());
+ if (update.getLastLogLineNr() != null && previous != null) {
+ String oldText = previous.getLoggingText();
+ if (StringUtils.isNotEmpty(oldText) && StringUtils.isNotEmpty(delta)) {
+ delta = capLoggingText(oldText + delta);
+ } else if (StringUtils.isNotEmpty(oldText)) {
+ delta = capLoggingText(oldText);
+ }
+ }
+ update.setLoggingText(delta);
+ }
+
+ private static String capLoggingText(String text) {
+ if (text == null || text.length() <= MAX_CACHED_LOGGING_TEXT_CHARS) {
+ return text;
+ }
+ return text.substring(text.length() - MAX_CACHED_LOGGING_TEXT_CHARS);
+ }
+
protected synchronized void addStateToCache(ExecutionState executionState)
throws HopException {
CacheEntry entry = cache.get(executionState.getId());
if (entry == null) {
@@ -397,6 +502,7 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
if (entry != null) {
// setExecutionState flags dirty so close()/timer flush include the state
+ appendLoggingDelta(entry.getExecutionState(), executionState);
entry.setExecutionState(executionState);
} else {
LogChannel.GENERAL.logError(
@@ -406,10 +512,11 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
}
- protected CacheEntry findCacheEntryWithParent(String parentId) {
+ protected synchronized CacheEntry findCacheEntryWithParent(String parentId) {
if (StringUtils.isEmpty(parentId)) {
return null;
}
+ CacheEntry found = null;
Collection<CacheEntry> values = cache.values();
for (CacheEntry cacheEntry : values) {
if (cacheEntry.getExecution() == null) {
@@ -417,15 +524,20 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
String execParent = cacheEntry.getExecution().getParentId();
if (parentId.equals(execParent)) {
- return cacheEntry;
- }
- if (cacheEntry.getExecutionState() == null) {
- continue;
+ found = cacheEntry;
+ break;
}
- if (parentId.equals(cacheEntry.getExecutionState().getId())) {
- return cacheEntry;
+ ExecutionState state = cacheEntry.peekExecutionState();
+ if (state != null && parentId.equals(state.getId())) {
+ found = cacheEntry;
+ break;
}
}
+ if (found != null) {
+ cache.get(found.getId());
+ found.markRead();
+ return found;
+ }
return null;
}
@@ -467,27 +579,39 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
protected synchronized CacheEntry findCacheEntry(String executionId) throws
HopException {
// Check the cache first...
+ CacheEntry found = null;
for (CacheEntry cacheEntry : cache.values()) {
// See if this is a parent in the cache.
//
if (cacheEntry.getId().equals(executionId)) {
- return cacheEntry;
+ found = cacheEntry;
+ break;
}
- // Sometimes the ID of the execution state is different from the
execution
+ // Sometimes the ID of the execution state is different from the
execution.
+ // Peek: a scan must not refresh lastRead on the entries that did not
match.
//
- if (cacheEntry.getExecutionState() != null
- && cacheEntry.getExecutionState().getId().equals(executionId)) {
- return cacheEntry;
+ ExecutionState state = cacheEntry.peekExecutionState();
+ if (state != null && executionId.equals(state.getId())) {
+ found = cacheEntry;
+ break;
}
// Is it perhaps one of the children?
//
- Execution childExecution = cacheEntry.getChildExecution(executionId);
+ Execution childExecution = cacheEntry.peekChildExecution(executionId);
if (childExecution != null) {
- return cacheEntry;
+ found = cacheEntry;
+ break;
}
}
+ if (found != null) {
+ // Iteration does not update access order. A key lookup moves this entry
to the newest end.
+ cache.get(found.getId());
+ found.markRead();
+ return found;
+ }
+
// We still haven't found anything in the cache.
// Let's load this from disk.
//
@@ -497,9 +621,8 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
entry.setLastWritten(new Date());
entry.setDirty(false);
- // Add this to the cache as well
- //
- cache.put(executionId, entry);
+ cache.put(entry.getId(), entry);
+ enforceMaxCacheSize();
return entry;
}
@@ -545,8 +668,11 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
CacheEntry entry = findCacheEntry(data.getParentId());
if (entry != null) {
entry.addExecutionData(data);
- // Flush promptly so other processes can merge samples when they persist
parent state
- persistCacheEntry(entry);
+ // A local single-writer flushes from the cache timer. Another process
(Spark/Beam) still
+ // needs the samples on disk before it merges its own write.
+ if (!entry.isSingleWriter()) {
+ persistCacheEntry(entry);
+ }
} else {
LogChannel.GENERAL.logError(
"Unable to register execution data for owner '"
@@ -593,7 +719,7 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
}
- protected void getExecutionIdsFromCache(Set<DatedId> ids, boolean
includeChildren) {
+ protected synchronized void getExecutionIdsFromCache(Set<DatedId> ids,
boolean includeChildren) {
for (CacheEntry cacheEntry : cache.values()) {
ids.add(new DatedId(cacheEntry.getId(),
cacheEntry.getExecution().getRegistrationDate()));
if (includeChildren) {
@@ -602,7 +728,8 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
}
- protected void getExecutionIdsFromCache(Set<DatedId> ids, IExecutionSelector
selector) {
+ protected synchronized void getExecutionIdsFromCache(
+ Set<DatedId> ids, IExecutionSelector selector) {
for (CacheEntry cacheEntry : cache.values()) {
if (selector.isSelected(cacheEntry.getExecution())
&& selector.isSelected(cacheEntry.getExecutionState())) {
@@ -615,12 +742,32 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
@Override
- public Execution getExecution(String executionId) throws HopException {
+ public synchronized Execution getExecution(String executionId) throws
HopException {
CacheEntry entry = findCacheEntry(executionId);
if (entry == null) {
return null;
}
- return entry.getExecution();
+ Execution execution = entry.getExecution();
+ if (execution != null
+ && entry.isHeavyDocumentStored()
+ && execution.getMetadataJson() == null
+ && execution.getExecutorXml() == null) {
+ restoreHeavyDocument(entry);
+ }
+ return execution;
+ }
+
+ /**
+ * The running entry drops the project metadata and pipeline XML after the
first insert. A viewer
+ * asking for the execution gets those fields back from the stored row, once.
+ */
+ private void restoreHeavyDocument(CacheEntry entry) throws HopException {
+ CacheEntry stored = loadCacheEntry(entry.getId());
+ if (stored == null || stored.getExecution() == null ||
entry.getExecution() == null) {
+ return;
+ }
+
entry.getExecution().setMetadataJson(stored.getExecution().getMetadataJson());
+
entry.getExecution().setExecutorXml(stored.getExecution().getExecutorXml());
}
@Override
diff --git
a/engine/src/main/java/org/apache/hop/execution/caching/CacheEntry.java
b/engine/src/main/java/org/apache/hop/execution/caching/CacheEntry.java
index 68bb027de8..c264fa3f56 100644
--- a/engine/src/main/java/org/apache/hop/execution/caching/CacheEntry.java
+++ b/engine/src/main/java/org/apache/hop/execution/caching/CacheEntry.java
@@ -82,6 +82,15 @@ public class CacheEntry {
// Was content modified and not yet written to disk?
@JsonIgnore private boolean dirty;
+ /**
+ * Local engine is the only writer of this row. Later saves update the small
columns and leave the
+ * stored document alone.
+ */
+ @JsonIgnore private boolean singleWriter;
+
+ /** The document with metadata and pipeline XML has been inserted. */
+ @JsonIgnore private boolean heavyDocumentStored;
+
public CacheEntry() {
childExecutions = new HashMap<>();
childExecutionStates = new HashMap<>();
@@ -194,6 +203,11 @@ public class CacheEntry {
return executionState;
}
+ /** Read the state without counting as a cache hit. Lookup scans must not
keep an entry warm. */
+ ExecutionState peekExecutionState() {
+ return executionState;
+ }
+
public void addChildExecution(Execution childExecution) {
childExecutions.put(childExecution.getId(), childExecution);
flagDirty();
@@ -216,6 +230,10 @@ public class CacheEntry {
}
}
+ void markRead() {
+ flagRead();
+ }
+
private void flagRead() {
lastRead = new Date();
}
@@ -225,6 +243,14 @@ public class CacheEntry {
return childExecutions.get(id);
}
+ /** Read a child without counting as a cache hit. */
+ Execution peekChildExecution(String id) {
+ if (childExecutions == null) {
+ return null;
+ }
+ return childExecutions.get(id);
+ }
+
public void addChildExecutionState(ExecutionState executionState) {
String executionId = executionState.getId();
if (StringUtils.isEmpty(executionId)) {
@@ -265,10 +291,14 @@ public class CacheEntry {
* @return true if this entry is too old.
*/
public boolean isTooOld(int maxAge) {
- if (lastRead != null && System.currentTimeMillis() - lastRead.getTime() >
maxAge) {
- return true;
+ long lastActivity = creationDate != null ? creationDate.getTime() : 0L;
+ if (lastWritten != null) {
+ lastActivity = Math.max(lastActivity, lastWritten.getTime());
+ }
+ if (lastRead != null) {
+ lastActivity = Math.max(lastActivity, lastRead.getTime());
}
- return lastWritten != null && System.currentTimeMillis() -
lastWritten.getTime() > maxAge;
+ return (System.currentTimeMillis() - lastActivity) > maxAge;
}
/**
diff --git
a/engine/src/main/java/org/apache/hop/execution/sampler/ExecutionDataSamplerStoreBase.java
b/engine/src/main/java/org/apache/hop/execution/sampler/ExecutionDataSamplerStoreBase.java
index ac8ed1d00e..0d1c2029b3 100644
---
a/engine/src/main/java/org/apache/hop/execution/sampler/ExecutionDataSamplerStoreBase.java
+++
b/engine/src/main/java/org/apache/hop/execution/sampler/ExecutionDataSamplerStoreBase.java
@@ -50,6 +50,13 @@ public abstract class ExecutionDataSamplerStoreBase<Store
extends IExecutionData
rows = Collections.synchronizedList(new ArrayList<>());
}
+ @Override
+ public void clearSamples() {
+ if (rows != null) {
+ rows.clear();
+ }
+ }
+
@Override
public IRowListener createRowListener(IExecutionDataSampler sampler) {
return new IRowListener() {
diff --git
a/engine/src/main/java/org/apache/hop/execution/sampler/IExecutionDataSamplerStore.java
b/engine/src/main/java/org/apache/hop/execution/sampler/IExecutionDataSamplerStore.java
index 3ed3af8ca1..79fc4a3716 100644
---
a/engine/src/main/java/org/apache/hop/execution/sampler/IExecutionDataSamplerStore.java
+++
b/engine/src/main/java/org/apache/hop/execution/sampler/IExecutionDataSamplerStore.java
@@ -57,4 +57,7 @@ public interface IExecutionDataSamplerStore {
* sampled data.
*/
Map<String, ExecutionDataSetMeta> getSamplesMetadata();
+
+ /** Drop captured rows. The next samples are collected from new rows only. */
+ default void clearSamples() {}
}
diff --git
a/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerStore.java
b/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerStore.java
index 5b259a167d..628c80353c 100644
---
a/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerStore.java
+++
b/engine/src/main/java/org/apache/hop/execution/sampler/plugins/dataprof/BasicDataProfilingDataSamplerStore.java
@@ -109,6 +109,13 @@ public class BasicDataProfilingDataSamplerStore
setMaxRows(Const.toInt(variables.resolve(dataSampler.getSampleSize()), 0));
}
+ @Override
+ public void clearSamples() {
+ super.clearSamples();
+ // Row buffers grow with the stream. Min, max and counters stay so a long
run keeps its profile.
+ profileSamples.clear();
+ }
+
@Override
public Map<String, RowBuffer> getSamples() {
Map<String, RowBuffer> samples = Collections.synchronizedMap(new
HashMap<>());
diff --git a/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
b/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
index a01b0ee78a..ee4a3de3e7 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/Pipeline.java
@@ -1558,38 +1558,44 @@ public abstract class Pipeline
@Override
public void fireExecutionFinishedListeners() throws HopException {
+ HopException listenerException = null;
synchronized (executionFinishedListeners) {
- if (executionFinishedListeners.isEmpty()) {
- return;
- }
- // prevent Exception from one listener to block others execution
- List<HopException> badGuys = new
ArrayList<>(executionFinishedListeners.size());
- for (IExecutionFinishedListener<IPipelineEngine<PipelineMeta>> listener :
- executionFinishedListeners) {
- try {
- listener.finished(this);
- } catch (HopException e) {
- badGuys.add(e);
+ if (!executionFinishedListeners.isEmpty()) {
+ // prevent Exception from one listener to block others execution
+ List<HopException> badGuys = new
ArrayList<>(executionFinishedListeners.size());
+ for (IExecutionFinishedListener<IPipelineEngine<PipelineMeta>>
listener :
+ executionFinishedListeners) {
+ try {
+ listener.finished(this);
+ } catch (HopException e) {
+ badGuys.add(e);
+ }
+ }
+ if (!badGuys.isEmpty()) {
+ // FIFO
+ listenerException = badGuys.get(0);
}
- }
- if (!badGuys.isEmpty()) {
- // FIFO
- throw new HopException(badGuys.get(0));
}
}
- // Now the status and everything else is set correctly. We've completed
the pipeline.
- //
- pipelineCompleted();
+ try {
+ // Now the status and everything else is set correctly. We've completed
the pipeline.
+ //
+ pipelineCompleted();
- // Also call an extension point in case plugins want to play along
- //
- ExtensionPointHandler.callExtensionPoint(
- log, this, HopExtensionPoint.PipelineCompleted.id, this);
+ // Also call an extension point in case plugins want to play along
+ //
+ ExtensionPointHandler.callExtensionPoint(
+ log, this, HopExtensionPoint.PipelineCompleted.id, this);
+ } finally {
+ // Only now: everything above can still touch files of this namespace,
and closing it
+ // invalidates every file object resolved through it - the result files
carry those.
+ releaseVfsNamespace();
+ }
- // Only now: everything above can still touch files of this namespace, and
closing it
- // invalidates every file object resolved through it - the result files
carry those.
- releaseVfsNamespace();
+ if (listenerException != null) {
+ throw listenerException;
+ }
}
public void pipelineCompleted() throws HopException {
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/engines/loadbalance/LoadBalancingPipelineEngine.java
b/engine/src/main/java/org/apache/hop/pipeline/engines/loadbalance/LoadBalancingPipelineEngine.java
index c31dbc38e3..3d00f0e777 100644
---
a/engine/src/main/java/org/apache/hop/pipeline/engines/loadbalance/LoadBalancingPipelineEngine.java
+++
b/engine/src/main/java/org/apache/hop/pipeline/engines/loadbalance/LoadBalancingPipelineEngine.java
@@ -19,6 +19,7 @@ package org.apache.hop.pipeline.engines.loadbalance;
import java.util.HashMap;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.lang3.StringUtils;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.logging.LogChannel;
@@ -54,6 +55,7 @@ public class LoadBalancingPipelineEngine extends
RemotePipelineEngine {
private LoadBalancingCoordinator<LoadBalancingPipelineRunConfiguration>
coordinator;
private LoadBalancingAssignment assignment;
private ExecutionInfoLocation executionInfoLocation;
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
public LoadBalancingPipelineEngine() {
super();
@@ -206,6 +208,16 @@ public class LoadBalancingPipelineEngine extends
RemotePipelineEngine {
}
}
+ /** Lines since the previous update. A full snapshot would be appended to
lines already stored. */
+ private ExecutionState captureExecutionState() {
+ ExecutionState state =
+ ExecutionStateBuilder.fromExecutor(this,
executionInfoLastLogLineNr.get()).build();
+ if (state.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(state.getLastLogLineNr());
+ }
+ return state;
+ }
+
private void registerLoadBalancingExecutionInformation(
ServerHealthSnapshot snapshot, int attempt) {
try {
@@ -215,7 +227,7 @@ public class LoadBalancingPipelineEngine extends
RemotePipelineEngine {
}
IExecutionInfoLocation location =
executionInfoLocation.getExecutionInfoLocation();
location.registerExecution(ExecutionBuilder.fromExecutor(this).build());
- ExecutionState state = ExecutionStateBuilder.fromExecutor(this,
-1).build();
+ ExecutionState state = captureExecutionState();
Map<String, String> details =
state.getDetails() == null ? new HashMap<>() : state.getDetails();
details.put(DETAIL_ASSIGNED_SERVER, selectedHopServerName);
@@ -267,7 +279,7 @@ public class LoadBalancingPipelineEngine extends
RemotePipelineEngine {
}
if (executionInfoLocation != null) {
IExecutionInfoLocation location =
executionInfoLocation.getExecutionInfoLocation();
- ExecutionState state = ExecutionStateBuilder.fromExecutor(this,
-1).build();
+ ExecutionState state = captureExecutionState();
Map<String, String> details =
state.getDetails() == null ? new HashMap<>() : state.getDetails();
if (assignment != null) {
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
index 7d77daea84..01da5371ed 100644
---
a/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
+++
b/engine/src/main/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngine.java
@@ -24,6 +24,7 @@ import java.util.Map;
import java.util.Timer;
import java.util.TimerTask;
import java.util.UUID;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.lang3.StringUtils;
import org.apache.hop.core.Const;
import org.apache.hop.core.IExtensionData;
@@ -43,6 +44,8 @@ import org.apache.hop.execution.ExecutionInfoLocation;
import org.apache.hop.execution.ExecutionState;
import org.apache.hop.execution.ExecutionStateBuilder;
import org.apache.hop.execution.IExecutionInfoLocation;
+import org.apache.hop.execution.caching.BaseCachingExecutionInfoLocation;
+import org.apache.hop.execution.caching.CacheEntry;
import org.apache.hop.execution.profiling.ExecutionDataProfile;
import org.apache.hop.execution.sampler.ExecutionDataSamplerMeta;
import org.apache.hop.execution.sampler.IExecutionDataSampler;
@@ -70,8 +73,15 @@ public class LocalPipelineEngine extends Pipeline implements
IPipelineEngine<Pip
private PipelineEngineCapabilities engineCapabilities = new
LocalPipelineEngineCapabilities();
private ExecutionInfoLocation executionInfoLocation;
+
+ void setExecutionInfoLocation(ExecutionInfoLocation executionInfoLocation) {
+ this.executionInfoLocation = executionInfoLocation;
+ }
+
private Timer transformExecutionInfoTimer;
private TimerTask transformExecutionInfoTimerTask;
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
+ private boolean executionInfoStoppedListenerRegistered;
private Map<String, List<IExecutionDataSamplerStore>> samplerStoresMap;
@@ -360,13 +370,16 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
executionInfoLocation
.getExecutionInfoLocation()
.registerExecution(ExecutionBuilder.fromExecutor(this).build());
+ markSingleWriter(executionInfoLocation.getExecutionInfoLocation(),
getLogChannelId());
// Also register an execution node for every transform
//
- for (TransformMetaDataCombi c : getTransforms()) {
- executionInfoLocation
- .getExecutionInfoLocation()
- .registerExecution(ExecutionBuilder.fromTransform(this,
c.transform).build());
+ if (getTransforms() != null) {
+ for (TransformMetaDataCombi c : getTransforms()) {
+ executionInfoLocation
+ .getExecutionInfoLocation()
+ .registerExecution(ExecutionBuilder.fromTransform(this,
c.transform).build());
+ }
}
}
}
@@ -448,7 +461,12 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
// Start execution of the pipeline
//
- super.startThreads();
+ try {
+ super.startThreads();
+ } catch (Exception e) {
+ stopTransformExecutionInfoTimer();
+ throw e;
+ }
// Make sure to kill the timer when this pipeline is finished.
// We do this with an extension point.
@@ -459,11 +477,29 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
super.waitUntilFinished();
}
- public void startTransformExecutionInfoTimer() throws HopException {
+ @Override
+ public void stopAll() {
+ super.stopAll();
+ // safeStop() only cancels the timer. A single-threaded child (a mapping,
the Kafka consumer,
+ // a Beam or Spark worker) never reaches pipelineCompleted(), so the
location is closed here.
+ // A normal pipeline sets finished and the end date after its transforms
stop, and
+ // pipelineCompleted() stores that terminal state. Closing here would drop
it.
+ if (getPipelineType() == PipelineMeta.PipelineType.SingleThreaded) {
+ stopTransformExecutionInfoTimer();
+ }
+ }
+
+ public synchronized void startTransformExecutionInfoTimer() throws
HopException {
if (executionInfoLocation == null) {
return;
}
+ if (!executionInfoStoppedListenerRegistered) {
+ addExecutionStoppedListener(e -> cancelTransformExecutionInfoTimer());
+ executionInfoStoppedListenerRegistered = true;
+ }
+ cancelTransformExecutionInfoTimer();
+
final ExecutionDataProfile dataProfile;
// No data profile to work with
@@ -482,47 +518,45 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
long interval =
Const.toLong(resolve(executionInfoLocation.getDataLoggingInterval()), 5000L);
final IExecutionInfoLocation iLocation =
executionInfoLocation.getExecutionInfoLocation();
- //
- TimerTask transformExecutionInfoTimerTask =
+ transformExecutionInfoTimerTask =
new TimerTask() {
@Override
public void run() {
- // Sample rows and execution state are written independently so a
conversion error
- // on sampled data cannot skip the state update (and hide the real
exception).
- //
- if (dataProfile != null) {
+ // Hold the engine lock before the location lock.
stopTransformExecutionInfoTimer()
+ // takes them in that order, and the reverse deadlocks a tick
against shutdown.
+ synchronized (LocalPipelineEngine.this) {
+ if (executionInfoLocation == null) {
+ return;
+ }
+ // Sample rows and execution state are written independently so
a conversion error
+ // on sampled data cannot skip the state update (and hide the
real exception).
+ //
+ if (dataProfile != null) {
+ try {
+ ExecutionDataBuilder dataBuilder =
+ ExecutionDataBuilder.fromAllTransformData(
+ LocalPipelineEngine.this, samplerStoresMap, false);
+ iLocation.registerData(dataBuilder.build());
+ } catch (Exception e) {
+ log.logError(
+ "Warning: unable to register execution data at location "
+ + executionInfoLocation.getName()
+ + " (non-fatal)",
+ e);
+ }
+ }
+
try {
- ExecutionDataBuilder dataBuilder =
- ExecutionDataBuilder.fromAllTransformData(
- LocalPipelineEngine.this, samplerStoresMap, false);
- iLocation.registerData(dataBuilder.build());
+ writeExecutionInfoState(iLocation);
+ releasePublishedSamples();
} catch (Exception e) {
log.logError(
- "Warning: unable to register execution data at location "
+ "Warning: unable to register execution state at location "
+ executionInfoLocation.getName()
+ " (non-fatal)",
e);
}
}
-
- try {
- ExecutionState pipelineState =
- ExecutionStateBuilder.fromExecutor(LocalPipelineEngine.this,
-1).build();
- iLocation.updateExecutionState(pipelineState);
-
- for (IEngineComponent component : getComponents()) {
- ExecutionState transformState =
-
ExecutionStateBuilder.fromTransform(LocalPipelineEngine.this, component)
- .build();
- iLocation.updateExecutionState(transformState);
- }
- } catch (Exception e) {
- log.logError(
- "Warning: unable to register execution state at location "
- + executionInfoLocation.getName()
- + " (non-fatal)",
- e);
- }
}
};
@@ -532,6 +566,62 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
transformExecutionInfoTimer.schedule(transformExecutionInfoTimerTask,
delay, interval);
}
+ /** This local engine is the only writer, so later saves can leave the
stored document alone. */
+ private void markSingleWriter(IExecutionInfoLocation location, String
executionId) {
+ if (location instanceof BaseCachingExecutionInfoLocation caching) {
+ synchronized (caching) {
+ CacheEntry entry = caching.getCache().get(executionId);
+ if (entry != null) {
+ entry.setSingleWriter(true);
+ }
+ }
+ }
+ }
+
+ /**
+ * Sample rows live on the engine for the whole run. Drop the row buffers
after they have been
+ * copied onto the execution. The copy stays in the cache so the next save
can publish it.
+ */
+ private void releasePublishedSamples() {
+ if (samplerStoresMap == null) {
+ return;
+ }
+ for (List<IExecutionDataSamplerStore> stores : samplerStoresMap.values()) {
+ for (IExecutionDataSamplerStore store : stores) {
+ store.clearSamples();
+ }
+ }
+ }
+
+ /**
+ * Asks the location for the lines after the previous tick. The caching
location appends that
+ * delta and drops the oldest characters, so this state does not keep the
whole log buffer.
+ */
+ private void writeExecutionInfoState(IExecutionInfoLocation iLocation)
throws HopException {
+ ExecutionState pipelineState =
+ ExecutionStateBuilder.fromExecutor(this,
executionInfoLastLogLineNr.get()).build();
+ if (pipelineState.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(pipelineState.getLastLogLineNr());
+ }
+ iLocation.updateExecutionState(pipelineState);
+
+ for (IEngineComponent component : getComponents()) {
+ ExecutionState transformState =
ExecutionStateBuilder.fromTransform(this, component).build();
+ iLocation.updateExecutionState(transformState);
+ }
+ }
+
+ private synchronized void cancelTransformExecutionInfoTimer() {
+ if (transformExecutionInfoTimerTask != null) {
+ transformExecutionInfoTimerTask.cancel();
+ transformExecutionInfoTimerTask = null;
+ }
+ if (transformExecutionInfoTimer != null) {
+ ExecutorUtil.cleanup(transformExecutionInfoTimer);
+ transformExecutionInfoTimer = null;
+ }
+ }
+
/**
* This method looks up the execution information location specified in the
run configuration.
*
@@ -565,61 +655,46 @@ public class LocalPipelineEngine extends Pipeline
implements IPipelineEngine<Pip
super.pipelineCompleted();
}
- public void stopTransformExecutionInfoTimer() {
+ public synchronized void stopTransformExecutionInfoTimer() {
try {
- if (transformExecutionInfoTimer != null) {
- if (transformExecutionInfoTimerTask != null) {
- transformExecutionInfoTimerTask.cancel();
- }
- ExecutorUtil.cleanup(transformExecutionInfoTimer);
- transformExecutionInfoTimer = null;
- }
+ cancelTransformExecutionInfoTimer();
- if (executionInfoLocation == null) {
+ ExecutionInfoLocation location = executionInfoLocation;
+ executionInfoLocation = null;
+ if (location == null || location.getExecutionInfoLocation() == null) {
return;
}
- IExecutionInfoLocation iLocation =
executionInfoLocation.getExecutionInfoLocation();
-
- // Register one final last state of the pipeline
- //
- IPipelineEngine pipelineEngine = LocalPipelineEngine.this;
-
- ExecutionStateBuilder stateBuilder =
ExecutionStateBuilder.fromExecutor(pipelineEngine, -1);
- ExecutionState executionState = stateBuilder.build();
- iLocation.updateExecutionState(executionState);
-
- // Update the state of all the transforms one final time
- //
- for (IEngineComponent component : getComponents()) {
- ExecutionState transformState =
- ExecutionStateBuilder.fromTransform(LocalPipelineEngine.this,
component).build();
- iLocation.updateExecutionState(transformState);
- }
+ IExecutionInfoLocation iLocation = location.getExecutionInfoLocation();
- String dataProfileName =
resolve(pipelineRunConfiguration.getExecutionDataProfileName());
- if (StringUtils.isNotEmpty(dataProfileName)) {
- // Register the collected transform data for the last time
+ try {
+ // Register one final last state of the pipeline
+ //
+ writeExecutionInfoState(iLocation);
+
+ String dataProfileName =
resolve(pipelineRunConfiguration.getExecutionDataProfileName());
+ if (StringUtils.isNotEmpty(dataProfileName)) {
+ // Register the collected transform data for the last time
+ //
+ ExecutionDataBuilder dataBuilder =
+ ExecutionDataBuilder.fromAllTransformData(
+ LocalPipelineEngine.this, samplerStoresMap, true);
+ iLocation.registerData(dataBuilder.build());
+ releasePublishedSamples();
+ }
+ } catch (Throwable e) {
+ log.logError("Error handling writing final pipeline state to location
(non-fatal)", e);
+ } finally {
+ // We're now certain all listeners fired. We can close the location.
//
- ExecutionDataBuilder dataBuilder =
- ExecutionDataBuilder.fromAllTransformData(
- LocalPipelineEngine.this, samplerStoresMap, true);
- iLocation.registerData(dataBuilder.build());
- }
- } catch (Throwable e) {
- log.logError("Error handling writing final pipeline state to location
(non-fatal)", e);
- } finally {
- // We're now certain all listeners fired. We can close the location.
- //
- if (executionInfoLocation != null) {
try {
- executionInfoLocation.getExecutionInfoLocation().close();
+ iLocation.close();
} catch (Exception e) {
- log.logError(
- "Error closing execution information location: " +
executionInfoLocation.getName(),
- e);
+ log.logError("Error closing execution information location: " +
location.getName(), e);
}
}
+ } catch (Throwable e) {
+ log.logError("Error stopping transform execution info timer
(non-fatal)", e);
}
}
diff --git
a/engine/src/main/java/org/apache/hop/workflow/engines/loadbalance/LoadBalancingWorkflowEngine.java
b/engine/src/main/java/org/apache/hop/workflow/engines/loadbalance/LoadBalancingWorkflowEngine.java
index 3c01cc5d31..032b398b0e 100644
---
a/engine/src/main/java/org/apache/hop/workflow/engines/loadbalance/LoadBalancingWorkflowEngine.java
+++
b/engine/src/main/java/org/apache/hop/workflow/engines/loadbalance/LoadBalancingWorkflowEngine.java
@@ -19,6 +19,7 @@ package org.apache.hop.workflow.engines.loadbalance;
import java.util.HashMap;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.commons.lang3.StringUtils;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.execution.ExecutionBuilder;
@@ -52,6 +53,7 @@ public class LoadBalancingWorkflowEngine extends
RemoteWorkflowEngine {
private LoadBalancingCoordinator<LoadBalancingWorkflowRunConfiguration>
coordinator;
private LoadBalancingAssignment assignment;
private ExecutionInfoLocation executionInfoLocation;
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
@Override
public IWorkflowEngineRunConfiguration
createDefaultWorkflowEngineRunConfiguration() {
@@ -180,6 +182,16 @@ public class LoadBalancingWorkflowEngine extends
RemoteWorkflowEngine {
}
}
+ /** Lines since the previous update. A full snapshot would be appended to
lines already stored. */
+ private ExecutionState captureExecutionState() {
+ ExecutionState state =
+ ExecutionStateBuilder.fromExecutor(this,
executionInfoLastLogLineNr.get()).build();
+ if (state.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(state.getLastLogLineNr());
+ }
+ return state;
+ }
+
private void registerLoadBalancingExecutionInformation(
ServerHealthSnapshot snapshot, int attempt) {
try {
@@ -189,7 +201,7 @@ public class LoadBalancingWorkflowEngine extends
RemoteWorkflowEngine {
}
IExecutionInfoLocation location =
executionInfoLocation.getExecutionInfoLocation();
location.registerExecution(ExecutionBuilder.fromExecutor(this).build());
- ExecutionState state = ExecutionStateBuilder.fromExecutor(this,
-1).build();
+ ExecutionState state = captureExecutionState();
Map<String, String> details =
state.getDetails() == null ? new HashMap<>() : state.getDetails();
details.put(DETAIL_ASSIGNED_SERVER, selectedHopServerName);
@@ -242,7 +254,7 @@ public class LoadBalancingWorkflowEngine extends
RemoteWorkflowEngine {
}
if (executionInfoLocation != null) {
IExecutionInfoLocation location =
executionInfoLocation.getExecutionInfoLocation();
- ExecutionState state = ExecutionStateBuilder.fromExecutor(this,
-1).build();
+ ExecutionState state = captureExecutionState();
Map<String, String> details =
state.getDetails() == null ? new HashMap<>() : state.getDetails();
if (assignment != null) {
diff --git
a/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
b/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
index 951a67310f..84def4818c 100644
---
a/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
+++
b/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
@@ -68,6 +68,7 @@ public class LocalWorkflowEngine extends Workflow implements
IWorkflowEngine<Wor
private ExecutionInfoLocation executionInfoLocation;
private Timer executionInfoTimer;
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
public LocalWorkflowEngine() {
super();
@@ -386,7 +387,6 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
long delay =
Const.toLong(resolve(executionInfoLocation.getDataLoggingDelay()), 2000L);
long interval =
Const.toLong(resolve(executionInfoLocation.getDataLoggingInterval()), 5000L);
- final AtomicInteger lastLogLineNr = new AtomicInteger(0);
final IExecutionInfoLocation iLocation =
executionInfoLocation.getExecutionInfoLocation();
@@ -402,12 +402,13 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
// Update the workflow execution state regularly
//
ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(LocalWorkflowEngine.this,
lastLogLineNr.get())
+ ExecutionStateBuilder.fromExecutor(
+ LocalWorkflowEngine.this,
executionInfoLastLogLineNr.get())
.build();
rebindSparkTransformOwnerParent(executionState);
iLocation.updateExecutionState(executionState);
if (executionState.getLastLogLineNr() != null) {
- lastLogLineNr.set(executionState.getLastLogLineNr());
+
executionInfoLastLogLineNr.set(executionState.getLastLogLineNr());
}
} catch (Exception e) {
log.logError(
@@ -530,7 +531,12 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
// Register one final last state of the workflow
//
ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(LocalWorkflowEngine.this,
-1).build();
+ ExecutionStateBuilder.fromExecutor(
+ LocalWorkflowEngine.this, executionInfoLastLogLineNr.get())
+ .build();
+ if (executionState.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(executionState.getLastLogLineNr());
+ }
rebindSparkTransformOwnerParent(executionState);
iLocation.updateExecutionState(executionState);
} finally {
diff --git
a/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_en_US.properties
b/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_en_US.properties
index 4c81439b4f..228305f4af 100644
---
a/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_en_US.properties
+++
b/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_en_US.properties
@@ -19,8 +19,10 @@ CachingFileExecutionInfoLocation.RootFolder.Label = Folder
to store the files in
CachingFileExecutionInfoLocation.RootFolder.Tooltip = Select the folder in
which you'll find a single file per executed pipeline or workflow
CachingFileExecutionInfoLocation.PersistenceDelay.Label = The persistence
delay (in ms)
CachingFileExecutionInfoLocation.PersistenceDelay.Tooltip = The maximum time
to wait before persisting data in the cache to file
+CachingFileExecutionInfoLocation.MaxCacheSize.Label = Maximum cache size
+CachingFileExecutionInfoLocation.MaxCacheSize.Tooltip = The maximum number of
executions to keep in the in-memory cache (default: 50)
CachingFileExecutionInfoLocation.MaxCacheAge.Label = Maximum cache entry age
(in ms)
-CachingFileExecutionInfoLocation.MaxCacheAge.Tooltip = The maximum time to
keep a cache entry around without it being read or written to
+CachingFileExecutionInfoLocation.MaxCacheAge.Tooltip = The maximum time (in
ms) to keep a cache entry around without it being read or written to. A missing
value keeps 86400000 ms (1 day). New locations default to 600000 ms (10
minutes).
CachingFileExecutionInfoLocation.CreateParentFolder.Label = Create folder?
CachingFileExecutionInfoLocation.CreateParentFolder.Tooltip = Check this
option to create the specified folder if it doesn't exist.
diff --git
a/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_pt_BR.properties
b/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_pt_BR.properties
index c33fbe31d6..cc10dc27cb 100644
---
a/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_pt_BR.properties
+++
b/engine/src/main/resources/org/apache/hop/execution/caching/messages/messages_pt_BR.properties
@@ -21,7 +21,9 @@ CachingFileExecutionInfoLocation.RootFolder.Label=Pasta para
guardar os arquivos
CachingFileExecutionInfoLocation.RootFolder.Tooltip=Selecione a pasta em qual
voc\u00EA ir\u00E1 encontrar um \u00FAnico arquivo por pipeline ou workflow
executado
CachingFileExecutionInfoLocation.PersistenceDelay.Label=O tempo de atraso de
persist\u00EAncia (em milissegundos)
CachingFileExecutionInfoLocation.PersistenceDelay.Tooltip=O tempo m\u00E1ximo
para esperar antes de gravar dados no arquivo de cache
+CachingFileExecutionInfoLocation.MaxCacheSize.Label=Tamanho m\u00E1ximo do
cache
+CachingFileExecutionInfoLocation.MaxCacheSize.Tooltip=O n\u00FAmero
m\u00E1ximo de execu\u00E7\u00F5es a serem mantidas no cache de mem\u00F3ria
(padr\u00E3o: 50)
CachingFileExecutionInfoLocation.MaxCacheAge.Label=Idade m\u00E1xima para
entrada de cache (em milissegundos)
-CachingFileExecutionInfoLocation.MaxCacheAge.Tooltip=O tempo m\u00E1ximo para
manter uma entrada no cache quando ela n\u00E3o seja lida ou gravada
+CachingFileExecutionInfoLocation.MaxCacheAge.Tooltip=O tempo m\u00E1ximo (em
ms) para manter uma entrada no cache sem leitura ou grava\u00E7\u00E3o. Um
valor ausente mant\u00E9m 86400000 ms (1 dia). Novas localiza\u00E7\u00F5es
usam 600000 ms (10 minutos) por padr\u00E3o.
CachingFileExecutionInfoLocation.CreateParentFolder.Label=Criar diret\u00F3rio?
CachingFileExecutionInfoLocation.CreateParentFolder.Tooltip=Marque essa
op\u00E7\u00E3o para criar o diret\u00F3rio especificado se ele n\u00E3o
existir.
diff --git
a/engine/src/test/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocationTest.java
b/engine/src/test/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocationTest.java
new file mode 100644
index 0000000000..b6c62b4e41
--- /dev/null
+++
b/engine/src/test/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocationTest.java
@@ -0,0 +1,231 @@
+/*
+ * 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.execution.caching;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.Date;
+import java.util.HashSet;
+import java.util.Set;
+import java.util.UUID;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.logging.HopLogStore;
+import org.apache.hop.core.variables.Variables;
+import org.apache.hop.execution.Execution;
+import org.apache.hop.execution.ExecutionState;
+import org.apache.hop.execution.ExecutionType;
+import org.apache.hop.execution.IExecutionSelector;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+class BaseCachingExecutionInfoLocationTest {
+
+ @BeforeAll
+ static void initLogging() {
+ if (!HopLogStore.isInitialized()) {
+ HopLogStore.init();
+ }
+ }
+
+ @Test
+ void missingMaxCacheAgeKeepsTheOneDayDefault() throws Exception {
+ FakeLocation location = new FakeLocation();
+ assertEquals(BaseCachingExecutionInfoLocation.LEGACY_MAX_CACHE_AGE,
location.getMaxCacheAge());
+
+ location.initialize(new Variables(), null);
+ try {
+ assertEquals(86_400_000, location.maxAge);
+ } finally {
+ location.close();
+ }
+ }
+
+ @Test
+ void explicitMaxCacheAgeIsUsed() throws Exception {
+ FakeLocation location = new FakeLocation();
+
location.setMaxCacheAge(BaseCachingExecutionInfoLocation.NEW_LOCATION_MAX_CACHE_AGE);
+
+ location.initialize(new Variables(), null);
+ try {
+ assertEquals(600_000, location.maxAge);
+ } finally {
+ location.close();
+ }
+ }
+
+ @Test
+ void loggingTextAppendsDeltasAndKeepsTheNewestCharacters() throws Exception {
+ FakeLocation location = new FakeLocation();
+ String id = UUID.randomUUID().toString();
+ location.registerExecution(pipeline(id, "Pipeline"));
+
+ ExecutionState first = pipelineState(id, "hello", 5);
+ location.updateExecutionState(first);
+ assertEquals("hello", location.getExecutionState(id).getLoggingText());
+
+ location.updateExecutionState(pipelineState(id, " world", 9));
+ assertEquals("hello world",
location.getExecutionState(id).getLoggingText());
+
+ location.updateExecutionState(pipelineState(id, "replaced", null));
+ assertEquals("replaced", location.getExecutionState(id).getLoggingText());
+
+ String head =
"H".repeat(BaseCachingExecutionInfoLocation.MAX_CACHED_LOGGING_TEXT_CHARS);
+ String tail = "T".repeat(32);
+ location.updateExecutionState(pipelineState(id, head, 10));
+ location.updateExecutionState(pipelineState(id, tail, 11));
+ String capped = location.getExecutionState(id).getLoggingText();
+
assertEquals(BaseCachingExecutionInfoLocation.MAX_CACHED_LOGGING_TEXT_CHARS,
capped.length());
+ assertTrue(capped.endsWith(tail));
+ assertTrue(
+ capped.startsWith(
+ "H"
+ .repeat(
+
BaseCachingExecutionInfoLocation.MAX_CACHED_LOGGING_TEXT_CHARS
+ - tail.length())));
+ }
+
+ @Test
+ void lookupMissDoesNotRefreshLastRead() throws Exception {
+ FakeLocation location = new FakeLocation();
+ String id = UUID.randomUUID().toString();
+ location.registerExecution(pipeline(id, "Kept"));
+ Date readAt = new Date(1_000L);
+ location.getCache().get(id).setLastRead(readAt);
+
+ assertNull(location.findCacheEntry("missing"));
+ assertEquals(readAt, location.getCache().get(id).getLastRead());
+
+ assertNotNull(location.findCacheEntry(id));
+ assertTrue(location.getCache().get(id).getLastRead().after(readAt));
+ }
+
+ @Test
+ void evictionKeepsADirtyEntryWhenPersistFails() throws Exception {
+ FakeLocation location = new FakeLocation();
+ location.maxSize = 1;
+ String kept = UUID.randomUUID().toString();
+ String added = UUID.randomUUID().toString();
+ location.registerExecution(pipeline(kept, "Kept"));
+ location.getCache().get(kept).setDirty(true);
+ location.failPersist = true;
+
+ assertThrows(HopException.class, () ->
location.registerExecution(pipeline(added, "Added")));
+
+ assertEquals(2, location.getCache().size());
+ assertTrue(location.getCache().containsKey(kept));
+ assertTrue(location.getCache().containsKey(added));
+ }
+
+ @Test
+ void closeLeavesUnpersistedEntriesAndRetriesThem() throws Exception {
+ FakeLocation location = new FakeLocation();
+ String id = UUID.randomUUID().toString();
+ location.registerExecution(pipeline(id, "Retry"));
+ location.getCache().get(id).setDirty(true);
+ location.failPersist = true;
+
+ assertThrows(HopException.class, location::close);
+ assertTrue(location.getCache().containsKey(id));
+
+ location.failPersist = false;
+ location.close();
+ assertTrue(location.getCache().isEmpty());
+ assertTrue(location.persisted.contains(id));
+ }
+
+ private static Execution pipeline(String id, String name) {
+ Execution execution = new Execution();
+ execution.setId(id);
+ execution.setName(name);
+ execution.setExecutionType(ExecutionType.Pipeline);
+ execution.setRegistrationDate(new Date());
+ return execution;
+ }
+
+ private static ExecutionState pipelineState(
+ String id, String loggingText, Integer lastLogLineNr) {
+ ExecutionState state = new ExecutionState();
+ state.setId(id);
+ state.setExecutionType(ExecutionType.Pipeline);
+ state.setLoggingText(loggingText);
+ state.setLastLogLineNr(lastLogLineNr);
+ return state;
+ }
+
+ private static final class FakeLocation extends
BaseCachingExecutionInfoLocation {
+ private boolean failPersist;
+ private final Set<String> persisted = new HashSet<>();
+ private String pluginId;
+ private String pluginName;
+
+ @Override
+ public String getPluginId() {
+ return pluginId;
+ }
+
+ @Override
+ public void setPluginId(String pluginId) {
+ this.pluginId = pluginId;
+ }
+
+ @Override
+ public String getPluginName() {
+ return pluginName;
+ }
+
+ @Override
+ public void setPluginName(String pluginName) {
+ this.pluginName = pluginName;
+ }
+
+ @Override
+ public BaseCachingExecutionInfoLocation clone() {
+ return new FakeLocation();
+ }
+
+ @Override
+ protected void persistCacheEntry(CacheEntry cacheEntry) throws
HopException {
+ if (failPersist) {
+ throw new HopException("persist failed for " + cacheEntry.getId());
+ }
+ persisted.add(cacheEntry.getId());
+ cacheEntry.setDirty(false);
+ cacheEntry.setLastWritten(new Date());
+ }
+
+ @Override
+ protected CacheEntry loadCacheEntry(String executionId) {
+ return null;
+ }
+
+ @Override
+ protected void deleteCacheEntry(CacheEntry cacheEntry) {
+ // Not used by these tests.
+ }
+
+ @Override
+ protected void retrieveIds(
+ boolean includeChildren, Set<DatedId> ids, int limit,
IExecutionSelector selector) {
+ // Not used by these tests.
+ }
+ }
+}
diff --git
a/engine/src/test/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngineExecutionIdTest.java
b/engine/src/test/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngineExecutionIdTest.java
new file mode 100644
index 0000000000..4f04ef28b5
--- /dev/null
+++
b/engine/src/test/java/org/apache/hop/pipeline/engines/local/LocalPipelineEngineExecutionIdTest.java
@@ -0,0 +1,172 @@
+/*
+ * 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.pipeline.engines.local;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import org.apache.hop.core.logging.HopLogStore;
+import org.apache.hop.core.logging.ILogChannel;
+import org.apache.hop.core.logging.LogChannel;
+import org.apache.hop.execution.Execution;
+import org.apache.hop.execution.ExecutionInfoLocation;
+import org.apache.hop.execution.ExecutionState;
+import org.apache.hop.execution.ExecutionStateBuilder;
+import org.apache.hop.execution.IExecutionSelector;
+import org.apache.hop.execution.caching.BaseCachingExecutionInfoLocation;
+import org.apache.hop.execution.caching.CacheEntry;
+import org.apache.hop.execution.caching.DatedId;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+
+class LocalPipelineEngineExecutionIdTest {
+
+ @BeforeAll
+ static void initLogging() {
+ if (!HopLogStore.isInitialized()) {
+ HopLogStore.init();
+ }
+ }
+
+ @Test
+ void registeredExecutionUsesTheLogChannelTheTimerWillRead() throws Exception
{
+ ILogChannel channel = new LogChannel("kafka-consumer");
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ pipelineMeta.setName("kafka-sub");
+ LocalPipelineEngine engine = new LocalPipelineEngine(pipelineMeta);
+ engine.setLogChannel(channel);
+ engine.setMetadataProvider(
+ new
org.apache.hop.metadata.serializer.memory.MemoryMetadataProvider());
+
+ CapturingLocation plugin = new CapturingLocation();
+ ExecutionInfoLocation location = new ExecutionInfoLocation();
+ location.setExecutionInfoLocation(plugin);
+ engine.setExecutionInfoLocation(location);
+
+ engine.registerPipelineExecutionInformation();
+
+ assertEquals(channel.getLogChannelId(), plugin.registeredId);
+ assertEquals(
+ channel.getLogChannelId(), ExecutionStateBuilder.fromExecutor(engine,
0).build().getId());
+
assertTrue(plugin.getCache().get(channel.getLogChannelId()).isSingleWriter());
+ }
+
+ @Test
+ void aNegativeLineRequestIsAFullSnapshot() {
+ ILogChannel channel = new LogChannel("snapshot");
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ pipelineMeta.setName("snapshot");
+ LocalPipelineEngine engine = new LocalPipelineEngine(pipelineMeta);
+ engine.setLogChannel(channel);
+
+ ExecutionState full = ExecutionStateBuilder.fromExecutor(engine,
-1).build();
+ assertNull(full.getLastLogLineNr());
+ assertNull(ExecutionStateBuilder.fromExecutor(engine,
null).build().getLastLogLineNr());
+
+ ExecutionState delta = ExecutionStateBuilder.fromExecutor(engine,
0).build();
+ assertNotNull(delta.getLastLogLineNr());
+ }
+
+ @Test
+ void stopAllClosesTheLocationOnlyForASingleThreadedEngine() throws Exception
{
+ CapturingLocation normalLocation = new CapturingLocation();
+ LocalPipelineEngine normal = engineWith(normalLocation);
+ normal.stopAll();
+ assertEquals(0, normalLocation.closes);
+
+ CapturingLocation singleThreadedLocation = new CapturingLocation();
+ LocalPipelineEngine singleThreaded = engineWith(singleThreadedLocation);
+ singleThreaded.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
+ singleThreaded.stopAll();
+ assertEquals(1, singleThreadedLocation.closes);
+ }
+
+ private static LocalPipelineEngine engineWith(CapturingLocation plugin) {
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ pipelineMeta.setName("stop");
+ LocalPipelineEngine engine = new LocalPipelineEngine(pipelineMeta);
+ engine.setMetadataProvider(
+ new
org.apache.hop.metadata.serializer.memory.MemoryMetadataProvider());
+ ExecutionInfoLocation location = new ExecutionInfoLocation();
+ location.setExecutionInfoLocation(plugin);
+ engine.setExecutionInfoLocation(location);
+ return engine;
+ }
+
+ private static final class CapturingLocation extends
BaseCachingExecutionInfoLocation {
+ private String registeredId;
+ private int closes;
+
+ @Override
+ public synchronized void close() throws
org.apache.hop.core.exception.HopException {
+ closes++;
+ super.close();
+ }
+
+ @Override
+ public void registerExecution(Execution execution)
+ throws org.apache.hop.core.exception.HopException {
+ registeredId = execution.getId();
+ super.registerExecution(execution);
+ }
+
+ @Override
+ public BaseCachingExecutionInfoLocation clone() {
+ return new CapturingLocation();
+ }
+
+ @Override
+ protected void persistCacheEntry(CacheEntry cacheEntry) {
+ cacheEntry.setDirty(false);
+ }
+
+ @Override
+ protected CacheEntry loadCacheEntry(String executionId) {
+ return null;
+ }
+
+ @Override
+ protected void deleteCacheEntry(CacheEntry cacheEntry) {}
+
+ @Override
+ protected void retrieveIds(
+ boolean includeChildren,
+ java.util.Set<DatedId> ids,
+ int limit,
+ IExecutionSelector selector) {}
+
+ @Override
+ public String getPluginId() {
+ return "capturing";
+ }
+
+ @Override
+ public void setPluginId(String pluginId) {}
+
+ @Override
+ public String getPluginName() {
+ return "capturing";
+ }
+
+ @Override
+ public void setPluginName(String pluginName) {}
+ }
+}
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/BeamPipelineEngine.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/BeamPipelineEngine.java
index 689bc8f94e..95824d021e 100644
---
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/BeamPipelineEngine.java
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/BeamPipelineEngine.java
@@ -30,6 +30,7 @@ import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
import org.apache.beam.runners.core.metrics.DefaultMetricResults;
import org.apache.beam.runners.dataflow.DataflowRunner;
import org.apache.beam.runners.direct.DirectRunner;
@@ -180,6 +181,7 @@ public abstract class BeamPipelineEngine extends Variables
private ExecutionInfoLocation executionInfoLocation;
private Timer executionInfoTimer;
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
/** Plugins can use this to add additional data samplers to the pipeline. */
protected List<IExecutionDataSampler<? extends IExecutionDataSamplerStore>>
dataSamplers;
@@ -1182,9 +1184,21 @@ public abstract class BeamPipelineEngine extends
Variables
executionInfoTimer.schedule(sampleTask, delay, interval);
}
- protected void updatePipelineState(IExecutionInfoLocation iLocation) throws
HopException {
+ /**
+ * Lines written since the previous tick. A full snapshot ({@code -1}) would
be appended on top of
+ * the lines already stored.
+ */
+ protected ExecutionState capturePipelineExecutionState() {
ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(BeamPipelineEngine.this,
-1).build();
+ ExecutionStateBuilder.fromExecutor(this,
executionInfoLastLogLineNr.get()).build();
+ if (executionState.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(executionState.getLastLogLineNr());
+ }
+ return executionState;
+ }
+
+ protected void updatePipelineState(IExecutionInfoLocation iLocation) throws
HopException {
+ ExecutionState executionState = capturePipelineExecutionState();
iLocation.updateExecutionState(executionState);
// Also update the state of the components
@@ -1209,8 +1223,7 @@ public abstract class BeamPipelineEngine extends Variables
// Register one final last state of the pipeline
//
- ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(BeamPipelineEngine.this,
-1).build();
+ ExecutionState executionState = capturePipelineExecutionState();
executionInfoLocation.getExecutionInfoLocation().updateExecutionState(executionState);
// Also update the state of the components
diff --git
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/dataflow/BeamDataFlowPipelineEngine.java
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/dataflow/BeamDataFlowPipelineEngine.java
index 25bfa38966..ac18eaf6e0 100644
---
a/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/dataflow/BeamDataFlowPipelineEngine.java
+++
b/plugins/engines/beam/src/main/java/org/apache/hop/beam/engines/dataflow/BeamDataFlowPipelineEngine.java
@@ -63,8 +63,7 @@ public class BeamDataFlowPipelineEngine extends
BeamPipelineEngine
@Override
protected void updatePipelineState(IExecutionInfoLocation iLocation) throws
HopException {
- ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(BeamDataFlowPipelineEngine.this,
-1).build();
+ ExecutionState executionState = capturePipelineExecutionState();
// Add Dataflow specific information to the execution state.
// This can then be picked up
diff --git
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
index a410b89c62..47b23ed6ed 100644
---
a/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
+++
b/plugins/engines/spark/src/main/java/org/apache/hop/spark/engines/SparkPipelineEngine.java
@@ -29,6 +29,8 @@ import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+import lombok.AccessLevel;
import lombok.Getter;
import lombok.Setter;
import org.apache.commons.lang3.StringUtils;
@@ -181,6 +183,10 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
/** Execution information location from the run configuration (optional). */
private ExecutionInfoLocation executionInfoLocation;
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private final AtomicInteger executionInfoLastLogLineNr = new
AtomicInteger(0);
+
private Timer executionInfoTimer;
private volatile boolean executionInfoClosed;
@@ -627,12 +633,24 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
interval);
}
+ /**
+ * Lines written since the previous tick. A full snapshot ({@code -1}) would
be appended on top of
+ * the lines already stored.
+ */
+ private ExecutionState capturePipelineExecutionState() {
+ ExecutionState executionState =
+ ExecutionStateBuilder.fromExecutor(this,
executionInfoLastLogLineNr.get()).build();
+ if (executionState.getLastLogLineNr() != null) {
+ executionInfoLastLogLineNr.set(executionState.getLastLogLineNr());
+ }
+ return executionState;
+ }
+
protected void updatePipelineState(IExecutionInfoLocation iLocation) throws
HopException {
// Register sample rows collected on executors before updating
parent/transform state
registerSampleDataFromExecutors(iLocation);
- ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(SparkPipelineEngine.this,
-1).build();
+ ExecutionState executionState = capturePipelineExecutionState();
iLocation.updateExecutionState(executionState);
// Transform Execution + state nodes under the parent pipeline (Beam does
the same from workers;
@@ -789,8 +807,7 @@ public class SparkPipelineEngine extends Variables
implements IPipelineEngine<Pi
// Final sample flush from executors (after jobs complete, accumulator
is fully merged)
registerSampleDataFromExecutors(iLocation);
- ExecutionState executionState =
- ExecutionStateBuilder.fromExecutor(SparkPipelineEngine.this,
-1).build();
+ ExecutionState executionState = capturePipelineExecutionState();
iLocation.updateExecutionState(executionState);
for (IEngineComponent component : getComponents()) {
diff --git
a/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
b/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
index f790666222..d48c0c0e0b 100644
---
a/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
+++
b/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
@@ -40,6 +40,7 @@ import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.gui.plugin.GuiElementType;
import org.apache.hop.core.gui.plugin.GuiPlugin;
import org.apache.hop.core.gui.plugin.GuiWidgetElement;
+import org.apache.hop.core.json.HopJson;
import org.apache.hop.core.logging.LogChannel;
import org.apache.hop.core.logging.LoggingObject;
import org.apache.hop.core.row.IRowMeta;
@@ -68,8 +69,8 @@ import org.apache.hop.ui.hopgui.HopGui;
/**
* Caches execution information and persists each top-level pipeline/workflow
{@link CacheEntry} as
- * one row in a relational table: filter columns for efficient queries, plus a
CLOB/TEXT column with
- * the full JSON payload (same shape as Caching File / Elastic / OpenSearch).
+ * one row in a relational table: filter columns for queries, a JSON document
written once (project
+ * metadata and pipeline XML), and a state document that a local single-writer
updates afterwards.
*/
@GuiPlugin(description = "Caching Database execution information location GUI
elements")
@ExecutionInfoLocationPlugin(
@@ -84,6 +85,8 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
public static final Class<?> PKG =
CachingDatabaseExecutionInfoLocation.class;
+ private static final ObjectMapper JSON_MAPPER = HopJson.newMapper();
+
public static final String PLUGIN_ID = "caching-database-location";
public static final String DEFAULT_TABLE_NAME = "hop_executions";
@@ -99,6 +102,12 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
public static final String COL_DURATION_MS = "duration_ms";
public static final String COL_JSON = "json";
+ /**
+ * Live state, child states and samples. {@link #COL_JSON} stays as inserted
so a long run does
+ * not bind the project metadata again.
+ */
+ public static final String COL_STATE_JSON = "state_json";
+
@GuiWidgetElement(
id = "connectionName",
order = "010",
@@ -149,6 +158,9 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
protected String actualSchemaName;
protected String actualTableName;
+ /** False only when an existing table could not grow the state column. New
DDL includes it. */
+ private boolean stateJsonColumnAvailable = true;
+
public CachingDatabaseExecutionInfoLocation() {
super();
}
@@ -212,6 +224,7 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
databaseClosed = false;
discardDatabase();
connectDatabase();
+ detectStateJsonColumn();
}
try {
@@ -242,6 +255,13 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
try {
super.close();
} finally {
+ // The caller closes once and drops this location. Nothing retries a
failed flush, so the
+ // connection has to go either way. Entries still marked dirty were not
saved.
+ if (hasDirtyCacheEntries()) {
+ LogChannel.GENERAL.logError(
+ "Closing the caching database execution information location with
unsaved entries for "
+ + getQuotedSchemaTable());
+ }
synchronized (dbLock) {
databaseClosed = true;
discardDatabase();
@@ -249,6 +269,15 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
}
}
+ private boolean hasDirtyCacheEntries() {
+ for (CacheEntry entry : getCache().values()) {
+ if (entry != null && entry.isDirty()) {
+ return true;
+ }
+ }
+ return false;
+ }
+
/**
* Open an independent auto-commit connection. A workflow transaction
connection is closed when
* the workflow stops, which left the cache timer writing through that
closed connection.
@@ -399,20 +428,44 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
throw new HopException("Caching database execution information location
is closed");
}
try {
- mergeChildrenFromDatabase(cacheEntry);
cacheEntry.calculateSummary();
-
- ObjectMapper mapper = new ObjectMapper();
- String json = mapper.writeValueAsString(cacheEntry);
-
- IRowMeta rowMeta = createDataRowMeta();
- Object[] data = buildRowData(cacheEntry, json);
-
- callWithDatabase(
- () -> {
- upsertCacheEntry(rowMeta, data);
- return null;
- });
+ // A light update leaves the inserted document alone and writes the
state to its own column.
+ // A table that predates that column has nowhere else to keep the state,
so it keeps
+ // rewriting the whole document.
+ boolean lightUpdate =
+ cacheEntry.isSingleWriter()
+ && cacheEntry.isHeavyDocumentStored()
+ && stateJsonColumnAvailable;
+ if (lightUpdate) {
+ callWithDatabase(
+ () -> {
+ if (rowExists(cacheEntry.getId())) {
+ updateLightState(cacheEntry);
+ return null;
+ }
+ // The row was removed. Write the document we still have.
+ writeFullDocument(cacheEntry);
+ return null;
+ });
+ } else {
+ // Another process may have added children. A local single-writer
inserts once and then
+ // never reads the CLOB back.
+ if (!cacheEntry.isSingleWriter()) {
+ mergeChildrenFromDatabase(cacheEntry);
+ }
+ callWithDatabase(
+ () -> {
+ writeFullDocument(cacheEntry);
+ return null;
+ });
+ // Dropping the metadata and the XML from memory is only safe once
later saves leave the
+ // stored document alone. Without the state column every save rewrites
it, so the entry
+ // has to keep them.
+ if (cacheEntry.isSingleWriter() && stateJsonColumnAvailable) {
+ releaseHeavyDocument(cacheEntry);
+ cacheEntry.setHeavyDocumentStored(true);
+ }
+ }
cacheEntry.setDirty(false);
cacheEntry.setLastWritten(new Date());
@@ -466,56 +519,125 @@ public class CachingDatabaseExecutionInfoLocation
extends BaseCachingExecutionIn
return one != null && one.getData() != null;
}
+ private void writeFullDocument(CacheEntry cacheEntry) throws HopException {
+ String json = serializeCacheEntry(cacheEntry);
+ IRowMeta rowMeta = createDataRowMeta();
+ // A full write is the source of truth, so drop any older state overlay.
+ Object[] data = buildRowData(cacheEntry, json, null);
+ upsertCacheEntry(rowMeta, data);
+ }
+
+ /**
+ * Status columns and the live state document. The inserted JSON, including
project metadata and
+ * pipeline XML, is left untouched and is not read back.
+ */
+ private void updateLightState(CacheEntry cacheEntry) throws HopException {
+ String stateJson = serializeWithoutHeavyDocument(cacheEntry);
+ Object[] data = buildRowData(cacheEntry, null, stateJson);
+ executeUpdate(lightUpdateFields(), data);
+ }
+
+ private String serializeWithoutHeavyDocument(CacheEntry cacheEntry) throws
HopException {
+ Execution execution = cacheEntry.getExecution();
+ String metadata = null;
+ String executorXml = null;
+ if (execution != null) {
+ metadata = execution.getMetadataJson();
+ executorXml = execution.getExecutorXml();
+ execution.setMetadataJson(null);
+ execution.setExecutorXml(null);
+ }
+ try {
+ return serializeCacheEntry(cacheEntry);
+ } finally {
+ if (execution != null) {
+ execution.setMetadataJson(metadata);
+ execution.setExecutorXml(executorXml);
+ }
+ }
+ }
+
+ private static String serializeCacheEntry(CacheEntry cacheEntry) throws
HopException {
+ try {
+ return JSON_MAPPER.writeValueAsString(cacheEntry);
+ } catch (Exception e) {
+ throw new HopException("Error serializing cache entry '" +
cacheEntry.getId() + "'", e);
+ }
+ }
+
+ private static String[] statusUpdateFields() {
+ return new String[] {
+ COL_NAME,
+ COL_EXECUTION_TYPE,
+ COL_PARENT_ID,
+ COL_REGISTRATION_DATE,
+ COL_EXECUTION_START_DATE,
+ COL_EXECUTION_END_DATE,
+ COL_FAILED,
+ COL_STATUS_DESCRIPTION,
+ COL_DURATION_MS
+ };
+ }
+
+ private String[] lightUpdateFields() {
+ String[] status = statusUpdateFields();
+ String[] fields = new String[status.length + 1];
+ System.arraycopy(status, 0, fields, 0, status.length);
+ fields[status.length] = COL_STATE_JSON;
+ return fields;
+ }
+
+ private String[] fullUpdateFields() {
+ String[] light = stateJsonColumnAvailable ? lightUpdateFields() :
statusUpdateFields();
+ String[] fields = new String[light.length + 1];
+ System.arraycopy(light, 0, fields, 0, light.length);
+ fields[light.length] = COL_JSON;
+ return fields;
+ }
+
+ /** prepareUpdate binds SET fields first, then the WHERE value. */
+ private void executeUpdate(String[] setFields, Object[] data) throws
HopException {
+ String[] codes = new String[] {COL_ID};
+ String[] conditions = new String[] {"="};
+ if (!database.prepareUpdate(actualSchemaName, actualTableName, codes,
conditions, setFields)) {
+ throw new HopException("Unable to prepare update for table " +
getQuotedSchemaTable());
+ }
+ try {
+ IRowMeta rowMeta = createDataRowMeta();
+ IRowMeta updateMeta = new RowMeta();
+ Object[] updateData = new Object[setFields.length + 1];
+ for (int i = 0; i < setFields.length; i++) {
+ int index = rowMeta.indexOfValue(setFields[i]);
+ if (index < 0) {
+ throw new HopException("Unknown execution information column '" +
setFields[i] + "'");
+ }
+ updateData[i] = data[index];
+ updateMeta.addValueMeta(rowMeta.getValueMeta(index));
+ }
+ int idIndex = rowMeta.indexOfValue(COL_ID);
+ updateData[setFields.length] = data[idIndex];
+ updateMeta.addValueMeta(rowMeta.getValueMeta(idIndex));
+ database.setValuesUpdate(updateMeta, updateData);
+ database.updateRow();
+ } finally {
+ database.closeUpdate();
+ }
+ }
+
+ private static void releaseHeavyDocument(CacheEntry cacheEntry) {
+ Execution execution = cacheEntry.getExecution();
+ if (execution == null) {
+ return;
+ }
+ execution.setMetadataJson(null);
+ execution.setExecutorXml(null);
+ }
+
/** Upsert: if a row with the same id exists UPDATE, otherwise INSERT. */
private void upsertCacheEntry(IRowMeta rowMeta, Object[] data) throws
HopException {
String id = (String) data[0];
if (rowExists(id)) {
- String[] setFields =
- new String[] {
- COL_NAME,
- COL_EXECUTION_TYPE,
- COL_PARENT_ID,
- COL_REGISTRATION_DATE,
- COL_EXECUTION_START_DATE,
- COL_EXECUTION_END_DATE,
- COL_FAILED,
- COL_STATUS_DESCRIPTION,
- COL_DURATION_MS,
- COL_JSON
- };
- String[] codes = new String[] {COL_ID};
- String[] conditions = new String[] {"="};
-
- if (!database.prepareUpdate(
- actualSchemaName, actualTableName, codes, conditions, setFields)) {
- throw new HopException("Unable to prepare update for table " +
getQuotedSchemaTable());
- }
- try {
- // prepareUpdate binds SET fields first, then WHERE values
- Object[] updateData = new Object[setFields.length + 1];
- updateData[0] = data[1];
- updateData[1] = data[2];
- updateData[2] = data[3];
- updateData[3] = data[4];
- updateData[4] = data[5];
- updateData[5] = data[6];
- updateData[6] = data[7];
- updateData[7] = data[8];
- updateData[8] = data[9];
- updateData[9] = data[10];
- updateData[10] = data[0];
-
- IRowMeta updateMeta = new RowMeta();
- for (String setField : setFields) {
- updateMeta.addValueMeta(rowMeta.searchValueMeta(setField));
- }
- updateMeta.addValueMeta(rowMeta.searchValueMeta(COL_ID));
-
- database.setValuesUpdate(updateMeta, updateData);
- database.updateRow();
- } finally {
- database.closeUpdate();
- }
+ executeUpdate(fullUpdateFields(), data);
} else {
database.insertRow(actualSchemaName, actualTableName, rowMeta, data);
}
@@ -524,9 +646,13 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
@Override
protected CacheEntry loadCacheEntry(String executionId) throws HopException {
try {
+ String columns = databaseMeta.quoteField(COL_JSON);
+ if (stateJsonColumnAvailable) {
+ columns += ", " + databaseMeta.quoteField(COL_STATE_JSON);
+ }
String sql =
"SELECT "
- + databaseMeta.quoteField(COL_JSON)
+ + columns
+ " FROM "
+ getQuotedSchemaTable()
+ " WHERE "
@@ -538,16 +664,18 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
return callWithDatabase(
() -> {
RowMetaAndData row = database.getOneRow(sql, paramMeta, new
Object[] {executionId});
- if (row == null || row.getData() == null) {
+ if (row == null || row.getData() == null || row.getData()[0] ==
null) {
return null;
}
- Object jsonObj = row.getData()[0];
- if (jsonObj == null) {
- return null;
+ CacheEntry entry =
JSON_MAPPER.readValue(row.getData()[0].toString(), CacheEntry.class);
+ if (stateJsonColumnAvailable
+ && row.getData().length > 1
+ && row.getData()[1] != null
+ && StringUtils.isNotEmpty(row.getData()[1].toString())) {
+ applyStateOverlay(
+ entry, JSON_MAPPER.readValue(row.getData()[1].toString(),
CacheEntry.class));
}
- String json = jsonObj.toString();
- ObjectMapper mapper = new ObjectMapper();
- return mapper.readValue(json, CacheEntry.class);
+ return entry;
});
} catch (Exception e) {
throw new HopException(
@@ -617,6 +745,7 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
callWithDatabase(
() -> {
+ List<DatedId> parentIds = new ArrayList<>();
ResultSet rs = database.openQuery(sql.toString(), paramMeta,
params.toArray());
try {
Object[] row = database.getRow(rs);
@@ -630,19 +759,23 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
startDate = (Date) row[1];
}
if (id != null) {
- ids.add(new DatedId(id, startDate != null ? startDate : new
Date(0L)));
- if (includeChildren && !activeSelector.isSelectingParents())
{
- CacheEntry entry = loadCacheEntry(id);
- if (entry != null) {
- addChildIds(entry, ids, activeSelector);
- }
- }
+ parentIds.add(new DatedId(id, startDate != null ? startDate
: new Date(0L)));
}
row = database.getRow(rs);
}
} finally {
database.closeQuery(rs);
}
+
+ ids.addAll(parentIds);
+ if (includeChildren && !activeSelector.isSelectingParents()) {
+ for (DatedId parentDatedId : parentIds) {
+ CacheEntry entry = loadCacheEntry(parentDatedId.getId());
+ if (entry != null) {
+ addChildIds(entry, ids, activeSelector);
+ }
+ }
+ }
return null;
});
} catch (Exception e) {
@@ -736,12 +869,15 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
rowMeta.addValueMeta(new ValueMetaString(COL_STATUS_DESCRIPTION, 128, -1));
// length 15 → BIGINT on most dialects (default Integer length maps to
tinyint on H2)
rowMeta.addValueMeta(new ValueMetaInteger(COL_DURATION_MS, 15, 0));
- // CLOB for full CacheEntry JSON
+ // CLOB for the CacheEntry JSON written on insert
rowMeta.addValueMeta(new ValueMetaString(COL_JSON,
DatabaseMeta.CLOB_LENGTH, -1));
+ if (stateJsonColumnAvailable) {
+ rowMeta.addValueMeta(new ValueMetaString(COL_STATE_JSON,
DatabaseMeta.CLOB_LENGTH, -1));
+ }
return rowMeta;
}
- private Object[] buildRowData(CacheEntry cacheEntry, String json) {
+ private Object[] buildRowData(CacheEntry cacheEntry, String json, String
stateJson) {
Execution execution = cacheEntry.getExecution();
ExecutionState state = cacheEntry.getExecutionState();
@@ -762,6 +898,21 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
Long durationMs =
cacheEntry.getSummary() != null ?
cacheEntry.getSummary().getDurationMs() : null;
+ if (!stateJsonColumnAvailable) {
+ return new Object[] {
+ cacheEntry.getId(),
+ name,
+ executionType,
+ parentId,
+ registrationDate,
+ startDate,
+ endDate,
+ failed,
+ status,
+ durationMs,
+ json
+ };
+ }
return new Object[] {
cacheEntry.getId(),
name,
@@ -773,10 +924,77 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
failed,
status,
durationMs,
- json
+ json,
+ stateJson
};
}
+ /**
+ * The state document is a later snapshot of the same cache entry without
project metadata or
+ * pipeline XML. Copy the parts that move onto the inserted document.
+ */
+ private static void applyStateOverlay(CacheEntry target, CacheEntry state) {
+ if (target == null || state == null) {
+ return;
+ }
+ if (state.getExecutionState() != null) {
+ target.setExecutionState(state.getExecutionState());
+ }
+ if (state.getChildExecutionStates() != null) {
+ target.setChildExecutionStates(state.getChildExecutionStates());
+ }
+ if (state.getChildExecutionData() != null) {
+ target.setChildExecutionData(state.getChildExecutionData());
+ }
+ if (state.getSummary() != null) {
+ target.setSummary(state.getSummary());
+ }
+ target.setDirty(false);
+ }
+
+ /**
+ * The state column is part of the DDL from {@link #buildDdl}. An existing
table that predates it
+ * keeps working through {@link #stateJsonColumnAvailable}. This location
does not alter the
+ * schema on startup.
+ */
+ private void detectStateJsonColumn() {
+ if (database == null || databaseMeta == null) {
+ return;
+ }
+ try {
+ if (!database.checkTableExists(actualSchemaName, actualTableName)) {
+ return;
+ }
+ stateJsonColumnAvailable =
+ database.checkColumnExists(actualSchemaName, actualTableName,
COL_STATE_JSON);
+ if (!stateJsonColumnAvailable) {
+ LogChannel.GENERAL.logBasic(
+ "Execution information table "
+ + getQuotedSchemaTable()
+ + " has no "
+ + COL_STATE_JSON
+ + " column, so every save rewrites the whole execution
document. Run this to store"
+ + " the execution state separately: "
+ + databaseMeta.getAddColumnStatement(
+ getQuotedSchemaTable(),
+ new ValueMetaString(COL_STATE_JSON,
DatabaseMeta.CLOB_LENGTH, -1),
+ "",
+ false,
+ "",
+ false));
+ }
+ } catch (Exception e) {
+ stateJsonColumnAvailable = false;
+ LogChannel.GENERAL.logError(
+ "Unable to see whether column "
+ + COL_STATE_JSON
+ + " exists on "
+ + getQuotedSchemaTable()
+ + ". Later updates keep the status columns only and leave
execution state in the inserted document.",
+ e);
+ }
+ }
+
protected String getQuotedSchemaTable() {
if (database != null) {
return databaseMeta.getQuotedSchemaTableCombination(
diff --git
a/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
b/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
index 04031320f2..682e24258b 100644
---
a/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
+++
b/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
@@ -35,11 +35,15 @@ import java.util.Timer;
import java.util.TimerTask;
import java.util.UUID;
import org.apache.hop.core.HopClientEnvironment;
+import org.apache.hop.core.RowMetaAndData;
import org.apache.hop.core.database.Database;
import org.apache.hop.core.database.DatabaseMeta;
import org.apache.hop.core.database.DatabasePluginType;
import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.logging.LoggingObject;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
import org.apache.hop.core.variables.Variables;
import org.apache.hop.databases.h2.H2DatabaseMeta;
import org.apache.hop.execution.DefaultExecutionSelector;
@@ -373,6 +377,7 @@ class CachingDatabaseExecutionInfoLocationTest {
String ddl = location.buildDdl(variables);
assertTrue(ddl.toLowerCase().contains("create"));
assertTrue(ddl.contains("idx_hop_exec_start") ||
ddl.toLowerCase().contains("index"));
+
assertTrue(ddl.toLowerCase().contains(CachingDatabaseExecutionInfoLocation.COL_STATE_JSON));
assertTrue(
ddl.contains(CachingDatabaseExecutionInfoLocation.COL_JSON)
|| ddl.toLowerCase().contains("json")
@@ -382,6 +387,318 @@ class CachingDatabaseExecutionInfoLocationTest {
|| ddl.toLowerCase().contains("character"));
}
+ @Test
+ void lruCacheEvictionEnforcesMaxSize() throws Exception {
+ location.setMaxCacheSize("2");
+ location.initialize(variables, metadataProvider);
+
+ String id1 = UUID.randomUUID().toString();
+ String id2 = UUID.randomUUID().toString();
+ String id3 = UUID.randomUUID().toString();
+
+ Execution exec1 = new Execution();
+ exec1.setId(id1);
+ exec1.setName("Exec1");
+ exec1.setExecutionType(ExecutionType.Pipeline);
+ exec1.setExecutionStartDate(new Date());
+ exec1.setRegistrationDate(new Date());
+
+ Execution exec2 = new Execution();
+ exec2.setId(id2);
+ exec2.setName("Exec2");
+ exec2.setExecutionType(ExecutionType.Pipeline);
+ exec2.setExecutionStartDate(new Date());
+ exec2.setRegistrationDate(new Date());
+
+ Execution exec3 = new Execution();
+ exec3.setId(id3);
+ exec3.setName("Exec3");
+ exec3.setExecutionType(ExecutionType.Pipeline);
+ exec3.setExecutionStartDate(new Date());
+ exec3.setRegistrationDate(new Date());
+
+ location.registerExecution(exec1);
+ location.registerExecution(exec2);
+ assertEquals(2, location.getCache().size());
+ assertTrue(location.getCache().containsKey(id1));
+ assertTrue(location.getCache().containsKey(id2));
+
+ // Registering the 3rd execution should evict the oldest (id1)
+ location.registerExecution(exec3);
+ assertEquals(2, location.getCache().size());
+ assertFalse(location.getCache().containsKey(id1));
+ assertTrue(location.getCache().containsKey(id2));
+ assertTrue(location.getCache().containsKey(id3));
+
+ // Evicted entry was persisted and can still be retrieved
+ Execution loaded1 = location.getExecution(id1);
+ assertNotNull(loaded1);
+ assertEquals("Exec1", loaded1.getName());
+ }
+
+ @Test
+ void closeClearsCacheMap() throws Exception {
+ String id = UUID.randomUUID().toString();
+ Execution exec = new Execution();
+ exec.setId(id);
+ exec.setName("ToClose");
+ exec.setExecutionType(ExecutionType.Pipeline);
+ exec.setExecutionStartDate(new Date());
+ exec.setRegistrationDate(new Date());
+
+ location.registerExecution(exec);
+ assertFalse(location.getCache().isEmpty());
+
+ location.close();
+ assertTrue(location.getCache().isEmpty());
+ }
+
+ @Test
+ void retrieveIdsWithChildrenLoadsChildrenCorrectly() throws Exception {
+ String parentId = UUID.randomUUID().toString();
+ String childId = UUID.randomUUID().toString();
+
+ CacheEntry parent =
+ sampleEntry(parentId, "ParentPipeline", ExecutionType.Pipeline, false,
"Finished");
+
+ Execution child = new Execution();
+ child.setId(childId);
+ child.setParentId(parentId);
+ child.setName("ChildPipeline");
+ child.setExecutionType(ExecutionType.Pipeline);
+ child.setExecutionStartDate(new Date());
+ child.setRegistrationDate(new Date());
+
+ parent.addChildExecution(child);
+ location.persistCacheEntry(parent);
+
+ // Clear memory cache so retrieveIds loads from DB
+ location.clearCaches();
+
+ Set<DatedId> ids = new HashSet<>();
+ location.retrieveIds(true, ids, 100, IExecutionSelector.ALL);
+ assertEquals(2, ids.size());
+ Set<String> idStrings = new HashSet<>();
+ ids.forEach(d -> idStrings.add(d.getId()));
+ assertTrue(idStrings.contains(parentId));
+ assertTrue(idStrings.contains(childId));
+ }
+
+ @Test
+ void singleWriterUpdateDoesNotReloadOrRewriteTheDocument() throws Exception {
+ CountingLocation counting = new CountingLocation();
+ counting.setConnectionName("h2-exec");
+
counting.setTableName(CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ counting.setPersistenceDelay("60000");
+ counting.setMaxCacheAge("86400000");
+ counting.setDatabaseMeta(databaseMeta);
+ counting.initialize(variables, metadataProvider);
+ try {
+ String id = UUID.randomUUID().toString();
+ CacheEntry entry = sampleEntry(id, "Live", ExecutionType.Pipeline,
false, "Running");
+ entry.getExecution().setMetadataJson("{\"project\":true}");
+ entry.getExecution().setExecutorXml("<pipeline/>");
+ entry.setSingleWriter(true);
+
+ counting.persistCacheEntry(entry);
+ assertNull(entry.getExecution().getMetadataJson());
+ assertNull(entry.getExecution().getExecutorXml());
+ assertTrue(entry.isHeavyDocumentStored());
+ assertEquals(0, counting.loads);
+
+ entry.getExecutionState().setStatusDescription("StillRunning");
+ entry.getExecution().setMetadataJson("SHOULD-NOT-BE-WRITTEN");
+ counting.persistCacheEntry(entry);
+ assertEquals(0, counting.loads);
+ assertEquals(
+ "StillRunning",
+ readColumn(id,
CachingDatabaseExecutionInfoLocation.COL_STATUS_DESCRIPTION));
+
+ String storedDocument = readColumn(id,
CachingDatabaseExecutionInfoLocation.COL_JSON);
+ String storedState = readColumn(id,
CachingDatabaseExecutionInfoLocation.COL_STATE_JSON);
+ assertFalse(storedDocument.contains("SHOULD-NOT-BE-WRITTEN"));
+ assertFalse(storedDocument.contains("StillRunning"));
+ assertTrue(storedState.contains("StillRunning"));
+ assertFalse(storedState.contains("SHOULD-NOT-BE-WRITTEN"));
+ assertFalse(storedState.contains("<pipeline/>"));
+
+ CacheEntry loaded = counting.loadCacheEntry(id);
+ assertEquals("{\"project\":true}",
loaded.getExecution().getMetadataJson());
+ assertEquals("<pipeline/>", loaded.getExecution().getExecutorXml());
+ assertEquals("StillRunning",
loaded.getExecutionState().getStatusDescription());
+ } finally {
+ counting.close();
+ }
+ }
+
+ @Test
+ void anExistingTableWithoutTheStateColumnIsNotAltered() throws Exception {
+ String table =
+ databaseMeta.getQuotedSchemaTableCombination(
+ variables, null,
CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ try (Database db = new Database(new LoggingObject("drop-state"),
variables, databaseMeta)) {
+ db.connect();
+ String sql =
+ databaseMeta.getDropColumnStatement(
+ table,
+ new
ValueMetaString(CachingDatabaseExecutionInfoLocation.COL_STATE_JSON),
+ "",
+ false,
+ "",
+ false);
+ db.execStatement(sql);
+ }
+
+ CachingDatabaseExecutionInfoLocation migrated = new
CachingDatabaseExecutionInfoLocation();
+ migrated.setConnectionName("h2-exec");
+
migrated.setTableName(CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ migrated.setPersistenceDelay("60000");
+ migrated.setMaxCacheAge("86400000");
+ migrated.setDatabaseMeta(databaseMeta);
+ migrated.initialize(variables, metadataProvider);
+ try {
+ assertFalse(
+ migrated.database.checkColumnExists(
+ null,
+ CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME,
+ CachingDatabaseExecutionInfoLocation.COL_STATE_JSON));
+
+ String id = UUID.randomUUID().toString();
+ migrated.persistCacheEntry(
+ sampleEntry(id, "Legacy", ExecutionType.Pipeline, false, "Running"));
+ assertEquals(
+ "Running", readColumn(id,
CachingDatabaseExecutionInfoLocation.COL_STATUS_DESCRIPTION));
+ } finally {
+ migrated.close();
+ }
+ }
+
+ @Test
+ void aTableWithoutTheStateColumnKeepsRewritingTheDocument() throws Exception
{
+ String table =
+ databaseMeta.getQuotedSchemaTableCombination(
+ variables, null,
CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ try (Database db =
+ new Database(new LoggingObject("drop-state-rewrite"), variables,
databaseMeta)) {
+ db.connect();
+ db.execStatement(
+ databaseMeta.getDropColumnStatement(
+ table,
+ new
ValueMetaString(CachingDatabaseExecutionInfoLocation.COL_STATE_JSON),
+ "",
+ false,
+ "",
+ false));
+ }
+
+ CachingDatabaseExecutionInfoLocation legacy = new
CachingDatabaseExecutionInfoLocation();
+ legacy.setConnectionName("h2-exec");
+
legacy.setTableName(CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ legacy.setPersistenceDelay("60000");
+ legacy.setMaxCacheAge("86400000");
+ legacy.setDatabaseMeta(databaseMeta);
+ legacy.initialize(variables, metadataProvider);
+ try {
+ String id = UUID.randomUUID().toString();
+ CacheEntry entry = sampleEntry(id, "Legacy", ExecutionType.Pipeline,
false, "Running");
+ entry.getExecution().setMetadataJson("{\"project\":true}");
+ entry.getExecution().setExecutorXml("<pipeline/>");
+ entry.setSingleWriter(true);
+ legacy.persistCacheEntry(entry);
+
+ entry.getExecutionState().setStatusDescription("Finished");
+ entry.getExecutionState().setLoggingText("FINAL-LOG-LINE");
+ entry.setDirty(true);
+ legacy.persistCacheEntry(entry);
+
+ // The state has nowhere else to go, so the document has to carry it.
+ CacheEntry reloaded = legacy.loadCacheEntry(id);
+ assertEquals("Finished",
reloaded.getExecutionState().getStatusDescription());
+ assertEquals("FINAL-LOG-LINE",
reloaded.getExecutionState().getLoggingText());
+
+ // Rewriting the document must not drop what was only stored on the
first save.
+ assertEquals("{\"project\":true}",
reloaded.getExecution().getMetadataJson());
+ assertEquals("<pipeline/>", reloaded.getExecution().getExecutorXml());
+
+ assertEquals(
+ "Finished", readColumn(id,
CachingDatabaseExecutionInfoLocation.COL_STATUS_DESCRIPTION));
+ } finally {
+ legacy.close();
+ }
+ }
+
+ @Test
+ void closeDropsTheConnectionWhenTheFinalFlushFails() throws Exception {
+ FailingFlushLocation failing = new FailingFlushLocation();
+ failing.setConnectionName("h2-exec");
+
failing.setTableName(CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ failing.setPersistenceDelay("60000");
+ failing.setMaxCacheAge("86400000");
+ failing.setDatabaseMeta(databaseMeta);
+ failing.initialize(variables, metadataProvider);
+ try {
+ String id = UUID.randomUUID().toString();
+ Execution execution = new Execution();
+ execution.setId(id);
+ execution.setName("Flush");
+ execution.setExecutionType(ExecutionType.Pipeline);
+ execution.setRegistrationDate(new Date());
+ failing.registerExecution(execution);
+ assertNotNull(failing.getCache().get(id));
+ failing.getCache().get(id).setDirty(true);
+ failing.failPersist = true;
+
+ assertThrows(HopException.class, failing::close);
+ assertNull(failing.database);
+ } finally {
+ failing.failPersist = false;
+ failing.close();
+ }
+ }
+
+ private static final class FailingFlushLocation extends
CachingDatabaseExecutionInfoLocation {
+ private boolean failPersist;
+
+ @Override
+ protected void persistCacheEntry(CacheEntry cacheEntry) throws
HopException {
+ if (failPersist) {
+ throw new HopException("flush failed");
+ }
+ super.persistCacheEntry(cacheEntry);
+ }
+ }
+
+ private static final class CountingLocation extends
CachingDatabaseExecutionInfoLocation {
+ private int loads;
+
+ @Override
+ protected CacheEntry loadCacheEntry(String executionId) throws
HopException {
+ loads++;
+ return super.loadCacheEntry(executionId);
+ }
+ }
+
+ private String readColumn(String id, String column) throws Exception {
+ try (Database db = new Database(new LoggingObject("status-check"),
variables, databaseMeta)) {
+ db.connect();
+ String table =
+ databaseMeta.getQuotedSchemaTableCombination(
+ variables, null,
CachingDatabaseExecutionInfoLocation.DEFAULT_TABLE_NAME);
+ String sql =
+ "SELECT "
+ + databaseMeta.quoteField(column)
+ + " FROM "
+ + table
+ + " WHERE "
+ +
databaseMeta.quoteField(CachingDatabaseExecutionInfoLocation.COL_ID)
+ + " = ?";
+ IRowMeta params = new RowMeta();
+ params.addValueMeta(new ValueMetaString("id"));
+ RowMetaAndData row = db.getOneRow(sql, params, new Object[] {id});
+ return row.getString(0, null);
+ }
+ }
+
private static CacheEntry sampleEntry(
String id, String name, ExecutionType type, boolean failed, String
status) {
Execution execution = new Execution();
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 75945faa39..59cf04e850 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
@@ -147,6 +147,10 @@ public class KafkaConsumerInput
kafkaPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
kafkaPipeline.setParentPipeline(getPipeline());
kafkaPipeline.setPipelineRunConfiguration(runConfiguration);
+ // Register under the consumer log channel. prepareExecution() captures
the id, and the
+ // execution-info timer later reads it again. Swapping the channel
afterwards made every tick
+ // miss the entry and keep it warm through a parent-id fallback.
+ kafkaPipeline.setLogChannel(getLogChannel());
kafkaPipeline.prepareExecution();
kafkaPipeline.setLogLevel(getPipeline().getLogLevel());
kafkaPipeline.setPreviousResult(new Result());
@@ -198,7 +202,6 @@ public class KafkaConsumerInput
}
});
}
- kafkaPipeline.setLogChannel(getLogChannel());
kafkaPipeline.startThreads();
if (errorHandlingConditionIsSatisfied()) {
diff --git
a/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
b/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
index 43c7f3d4d0..a4f94b9321 100644
---
a/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
+++
b/plugins/transforms/mapping/src/main/java/org/apache/hop/pipeline/transforms/mapping/SimpleMapping.java
@@ -161,6 +161,16 @@ public class SimpleMapping extends
BaseTransform<SimpleMappingMeta, SimpleMappin
if (isSingleThreaded()) {
simpleMappingData.mappingPipeline.setPipelineType(PipelineMeta.PipelineType.SingleThreaded);
+ // This child is driven again on every parent iteration and does not
finish between them.
+ // A location here would register a second caching session (metadata
document, sampler rows,
+ // and a timer) for the whole parent run.
+ if (suppressExecutionInformation(simpleMappingData.mappingPipeline)) {
+ logDetailed(
+ BaseMessages.getString(
+ PKG,
+ "SimpleMapping.Log.IgnoringExecutionInformationLocation",
+
simpleMappingData.mappingPipeline.getPipelineRunConfiguration().getName()));
+ }
}
// Copy the parameters over...
@@ -247,6 +257,24 @@ public class SimpleMapping extends
BaseTransform<SimpleMappingMeta, SimpleMappin
&& getPipeline().getPipelineType() ==
PipelineMeta.PipelineType.SingleThreaded);
}
+ /**
+ * Drops the execution-information location from a copy of the child's run
configuration. The
+ * metadata object itself is left unchanged.
+ *
+ * @return true when a location was removed
+ */
+ static boolean suppressExecutionInformation(LocalPipelineEngine engine) {
+ PipelineRunConfiguration runConfig = engine.getPipelineRunConfiguration();
+ if (runConfig == null ||
StringUtils.isEmpty(runConfig.getExecutionInfoLocationName())) {
+ return false;
+ }
+ PipelineRunConfiguration copy = new PipelineRunConfiguration(runConfig);
+ copy.setExecutionInfoLocationName(null);
+ copy.setExecutionDataProfileName(null);
+ engine.setPipelineRunConfiguration(copy);
+ return true;
+ }
+
public static List<MappingInput> findMappingInputs(Pipeline mappingPipeline)
{
return MappingTransforms.findMappingInputs(mappingPipeline);
}
@@ -312,6 +340,10 @@ public class SimpleMapping extends
BaseTransform<SimpleMappingMeta, SimpleMappin
try {
if (data.executor != null) {
try {
+ // A single-threaded child has no transform threads, so it does not
finish on its own.
+ if (data.mappingPipeline != null &&
!data.mappingPipeline.isFinished()) {
+ data.mappingPipeline.stopAll();
+ }
data.executor.dispose();
} catch (Exception e) {
logError("Error calling dispose() on single threaded Simple Mapping
executor", e);
diff --git
a/plugins/transforms/mapping/src/main/resources/org/apache/hop/pipeline/transforms/mapping/messages/messages_en_US.properties
b/plugins/transforms/mapping/src/main/resources/org/apache/hop/pipeline/transforms/mapping/messages/messages_en_US.properties
index 521011085e..cfc2ddc9aa 100644
---
a/plugins/transforms/mapping/src/main/resources/org/apache/hop/pipeline/transforms/mapping/messages/messages_en_US.properties
+++
b/plugins/transforms/mapping/src/main/resources/org/apache/hop/pipeline/transforms/mapping/messages/messages_en_US.properties
@@ -22,6 +22,7 @@
SimpleMapping.Exception.MoreThanOneRowReceivedFromTransform=More than 1 row of d
SimpleMapping.Exception.SameFieldMappedTwice=The same input field [{0}] was
mapped twice (or more). At present, this is not allowed. Please use a
''Select Values'' transform to make a copy of the field first.
SimpleMapping.Exception.UnableToInitSingleThreadedPipeline=Unable to
initialize the single threaded pipeline engine because one or more transforms
failed to initialize.
SimpleMapping.Exception.UnableToPrepareExecutionOfMapping=Unable to prepare
(allocate, initialize) the execution of the mapping (sub-pipeline)
+SimpleMapping.Log.IgnoringExecutionInformationLocation=Not using the execution
information location on run configuration ''{0}'' for this single-threaded
mapping. The parent pipeline tracks the run.
SimpleMapping.Exception.UnableToReadRowFromTransform=Unable to read a single
row of data from transform ''{0}''.
SimpleMapping.Log.CouldNotFindMappingInputTransform=Couldn''t find
MappingInput transform in the mapping.
SimpleMapping.Log.CouldNotFindMappingInputTransform2=Couldn''t find
MappingOutput transform in the mapping.
diff --git
a/plugins/transforms/mapping/src/test/java/org/apache/hop/pipeline/transforms/mapping/SimpleMappingExecutionInfoTest.java
b/plugins/transforms/mapping/src/test/java/org/apache/hop/pipeline/transforms/mapping/SimpleMappingExecutionInfoTest.java
new file mode 100644
index 0000000000..339cf50b20
--- /dev/null
+++
b/plugins/transforms/mapping/src/test/java/org/apache/hop/pipeline/transforms/mapping/SimpleMappingExecutionInfoTest.java
@@ -0,0 +1,60 @@
+/*
+ * 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.pipeline.transforms.mapping;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import org.apache.hop.pipeline.config.PipelineRunConfiguration;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.junit.jupiter.api.Test;
+
+class SimpleMappingExecutionInfoTest {
+
+ @Test
+ void singleThreadedMappingDropsTheChildExecutionInformationLocation() {
+ LocalPipelineEngine engine = new LocalPipelineEngine();
+ PipelineRunConfiguration original = engine.getPipelineRunConfiguration();
+ original.setName("with-location");
+ original.setExecutionInfoLocationName("caching-db");
+ original.setExecutionDataProfileName("first-rows");
+
+ assertTrue(SimpleMapping.suppressExecutionInformation(engine));
+
+ PipelineRunConfiguration used = engine.getPipelineRunConfiguration();
+ assertNull(used.getExecutionInfoLocationName());
+ assertNull(used.getExecutionDataProfileName());
+ assertEquals("with-location", used.getName());
+ assertEquals("caching-db", original.getExecutionInfoLocationName());
+ assertEquals("first-rows", original.getExecutionDataProfileName());
+ assertNotSame(original, used);
+ }
+
+ @Test
+ void mappingWithoutALocationKeepsItsRunConfiguration() {
+ LocalPipelineEngine engine = new LocalPipelineEngine();
+ PipelineRunConfiguration original = engine.getPipelineRunConfiguration();
+
+ assertFalse(SimpleMapping.suppressExecutionInformation(engine));
+ assertSame(original, engine.getPipelineRunConfiguration());
+ }
+}
diff --git
a/ui/src/main/java/org/apache/hop/ui/execution/ExecutionInfoLocationEditor.java
b/ui/src/main/java/org/apache/hop/ui/execution/ExecutionInfoLocationEditor.java
index 2ed98185da..a40c84634c 100644
---
a/ui/src/main/java/org/apache/hop/ui/execution/ExecutionInfoLocationEditor.java
+++
b/ui/src/main/java/org/apache/hop/ui/execution/ExecutionInfoLocationEditor.java
@@ -29,6 +29,7 @@ import org.apache.hop.core.plugins.IPlugin;
import org.apache.hop.core.plugins.PluginRegistry;
import org.apache.hop.execution.ExecutionInfoLocation;
import org.apache.hop.execution.IExecutionInfoLocation;
+import org.apache.hop.execution.caching.BaseCachingExecutionInfoLocation;
import org.apache.hop.execution.plugin.ExecutionInfoLocationPluginType;
import org.apache.hop.i18n.BaseMessages;
import org.apache.hop.ui.core.PropsUi;
@@ -108,6 +109,12 @@ public class ExecutionInfoLocationEditor extends
MetadataEditor<ExecutionInfoLoc
location.setPluginId(plugin.getIds()[0]);
location.setPluginName(plugin.getName());
+ if (location instanceof BaseCachingExecutionInfoLocation
cachingLocation) {
+ // Fresh plugin instances are the new-location path. A loaded
location replaces this
+ // object in the map, so an omitted maxCacheAge on disk stays at one
day.
+ cachingLocation.setMaxCacheAge(
+ BaseCachingExecutionInfoLocation.NEW_LOCATION_MAX_CACHE_AGE);
+ }
metaMap.put(plugin.getName(), location);
} catch (Exception e) {