This is an automated email from the ASF dual-hosted git repository.
bamaer 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 7658c7549a Issue #2406 : Flush text file output every few seconds by
default (#8674)
7658c7549a is described below
commit 7658c7549a738611db5467014ed8a3ef07c437fd
Author: Matt Casters <[email protected]>
AuthorDate: Thu Oct 1 20:49:10 2026 +0200
Issue #2406 : Flush text file output every few seconds by default (#8674)
* Issue #2406 : Flush text file output every few seconds by default
* Issue #2406 : Close text files after an interval flush and let -1 disable
it
---
core/src/main/java/org/apache/hop/core/Const.java | 9 +-
.../pages/pipeline/transforms/textfileoutput.adoc | 2 +
.../modules/ROOT/pages/variables.adoc | 9 +-
.../hop/core/compress/CompressionOutputStream.java | 7 +
.../core/compress/CompressionOutputStreamTest.java | 21 ++
.../transforms/textfileoutput/TextFileOutput.java | 68 +++--
.../textfileoutput/TextFileOutputData.java | 46 +--
.../textfileoutput/TextFileOutputFlushTest.java | 309 +++++++++++++++++++++
8 files changed, 429 insertions(+), 42 deletions(-)
diff --git a/core/src/main/java/org/apache/hop/core/Const.java
b/core/src/main/java/org/apache/hop/core/Const.java
index e916818a48..04382353ae 100644
--- a/core/src/main/java/org/apache/hop/core/Const.java
+++ b/core/src/main/java/org/apache/hop/core/Const.java
@@ -811,13 +811,14 @@ public class Const {
public static final String HOP_FILE_OUTPUT_MAX_STREAM_COUNT =
"HOP_FILE_OUTPUT_MAX_STREAM_COUNT";
/**
- * This variable contains the number of milliseconds between flushes of all
open files in the Text
- * File Output transform.
+ * Milliseconds between flushes of all open files in the Text File Output
transform. {@code 0}
+ * selects the transform default of 5000. A negative value, for example
{@code -1}, disables the
+ * interval flush.
*/
@Variable(
- value = "0",
+ value = "5000",
description =
- "This project variable is used by the Text File Output transform. It
defines the max number of milliseconds between flushes of files opened by the
transform.")
+ "This project variable is used by the Text File Output transform. It
defines how many milliseconds to wait between flushes of files opened by the
transform. Output is buffered, so slow input stays invisible until the buffer
fills or the file is closed. The default is 5000 (5 seconds). A value of 0 uses
that default. Set a positive number of milliseconds to change the interval. A
negative value, for example -1, disables the interval flush.")
public static final String HOP_FILE_OUTPUT_MAX_STREAM_LIFE =
"HOP_FILE_OUTPUT_MAX_STREAM_LIFE";
/** Set this variable to Y to disable standard Hop logging to the console.
(stdout) */
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/textfileoutput.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/textfileoutput.adoc
index 6a14d2e6d5..45f66ce161 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/textfileoutput.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/textfileoutput.adoc
@@ -29,6 +29,8 @@ This is commonly used to generate Comma Separated Values (CSV
files) that can be
It is also possible to generate fixed width files by setting lengths on the
fields in the fields tab.
+Output is buffered. Open files are flushed on an interval so slow-arriving
rows show up without waiting for that buffer to fill. The interval is
xref:variables.adoc[HOP_FILE_OUTPUT_MAX_STREAM_LIFE], in milliseconds. The
default is 5 seconds. A value of 0 uses that default. Set a positive number of
milliseconds to change the interval. A negative value, for example -1, disables
the interval flush.
+
You can choose to use a
xref:metadata-types/static-schema-definition.adoc[Schema Definition] or to
define the required fields' layout manually. If you decide to define the fields
layout by using a xref:metadata-types/static-schema-definition.adoc[Schema
Definition], use the xref:pipeline/transforms/schemamapping.adoc[Schema
mapping] transform to adjust the incoming stream according to the chosen
xref:metadata-types/static-schema-definition.adoc[Schema Definition]
The `ignore manual fields` ignores any fields manually defined in the
transform's field layout, and only uses the layout specified in the
xref:metadata-types/static-schema-definition.adoc[Schema Definition].
diff --git a/docs/hop-user-manual/modules/ROOT/pages/variables.adoc
b/docs/hop-user-manual/modules/ROOT/pages/variables.adoc
index e8ac3c38ef..e0b6af0efb 100644
--- a/docs/hop-user-manual/modules/ROOT/pages/variables.adoc
+++ b/docs/hop-user-manual/modules/ROOT/pages/variables.adoc
@@ -258,8 +258,13 @@ Otherwise they are not.
|HOP_FILE_OUTPUT_MAX_STREAM_COUNT|1024|This project variable is used by the
Text File Output transform.
It defines the max number of simultaneously open files within the transform.
The transform will close/reopen files as necessary to insure the max is not
exceeded
-|HOP_FILE_OUTPUT_MAX_STREAM_LIFE|0|This project variable is used by the Text
File Output transform.
-It defines the max number of milliseconds between flushes of files opened by
the transform.
+|HOP_FILE_OUTPUT_MAX_STREAM_LIFE|5000|This project variable is used by the
Text File Output transform.
+It defines how many milliseconds to wait between flushes of files opened by
the transform.
+Output is buffered, so slow input stays invisible until the buffer fills or
the file is closed.
+The default is 5000 (5 seconds).
+A value of 0 uses that default.
+Set a positive number of milliseconds to change the interval.
+A negative value, for example -1, disables the interval flush.
|HOP_GLOBAL_LOG_VARIABLES_CLEAR_ON_EXPORT|N|Set this variable to N to preserve
global log variables defined in pipeline / workflow Properties -> Log panel.
Changing it to true will clear it when export pipeline / workflow.
|HOP_JSON_INPUT_INCLUDE_NULLS|Y|Name of te variable to set so that Nulls are
considered while parsing JSON files. If HOP_JSON_INPUT_INCLUDE_NULLS is "Y"
then nulls will be included otherwise they will not be included (default
behavior)
diff --git
a/engine/src/main/java/org/apache/hop/core/compress/CompressionOutputStream.java
b/engine/src/main/java/org/apache/hop/core/compress/CompressionOutputStream.java
index 4c15cd9470..3753c26327 100644
---
a/engine/src/main/java/org/apache/hop/core/compress/CompressionOutputStream.java
+++
b/engine/src/main/java/org/apache/hop/core/compress/CompressionOutputStream.java
@@ -48,6 +48,13 @@ public abstract class CompressionOutputStream extends
OutputStream {
delegate.close();
}
+ @Override
+ public void flush() throws IOException {
+ if (delegate != null) {
+ delegate.flush();
+ }
+ }
+
@Override
public void write(int b) throws IOException {
delegate.write(b);
diff --git
a/engine/src/test/java/org/apache/hop/core/compress/CompressionOutputStreamTest.java
b/engine/src/test/java/org/apache/hop/core/compress/CompressionOutputStreamTest.java
index def2e2e9ec..939d59c2cb 100644
---
a/engine/src/test/java/org/apache/hop/core/compress/CompressionOutputStreamTest.java
+++
b/engine/src/test/java/org/apache/hop/core/compress/CompressionOutputStreamTest.java
@@ -80,6 +80,18 @@ class CompressionOutputStreamTest {
outStream.write("Test".getBytes());
}
+ @Test
+ void testFlushDelegates() throws IOException {
+ ICompressionProvider provider = outStream.getCompressionProvider();
+ FlushCountStream out = new FlushCountStream();
+ outStream = new DummyCompressionOS(out, provider);
+ outStream.write("Test".getBytes());
+ assertEquals(0, out.flushes);
+ outStream.flush();
+ assertEquals(1, out.flushes);
+ assertEquals("Test", out.toString());
+ }
+
@Test
void testAddEntry() throws IOException {
ICompressionProvider provider = outStream.getCompressionProvider();
@@ -94,4 +106,13 @@ class CompressionOutputStreamTest {
super(out, provider);
}
}
+
+ private static class FlushCountStream extends ByteArrayOutputStream {
+ private int flushes;
+
+ @Override
+ public void flush() {
+ flushes++;
+ }
+ }
}
diff --git
a/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutput.java
b/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutput.java
index 80df8b5810..002f981aac 100644
---
a/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutput.java
+++
b/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutput.java
@@ -22,7 +22,6 @@ import java.io.IOException;
import java.io.OutputStream;
import java.io.UnsupportedEncodingException;
import java.util.ArrayList;
-import java.util.Date;
import java.util.List;
import org.apache.commons.vfs2.FileObject;
import org.apache.hop.core.Const;
@@ -124,7 +123,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
@SuppressWarnings("java:S2095") // the stream is owned by the transform and
closed in closeFile()
public void initFileStreamWriter(String filename) throws HopException {
- data.writer = null;
+ assignWriter(null, null);
try {
TextFileOutputData.FileStream fileStreams = null;
@@ -248,7 +247,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
data.fos = fileStreams.getFileOutputStream();
data.out = fileStreams.getCompressedOutputStream();
- data.writer = fileStreams.getBufferedOutputStream();
+ assignWriter(fileStreams.getBufferedOutputStream(), fileStreams);
} catch (HopException ke) {
throw ke;
} catch (Exception e) {
@@ -289,15 +288,20 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
return filename;
}
+ /**
+ * Milliseconds between flushes when {@link
Const#HOP_FILE_OUTPUT_MAX_STREAM_LIFE} is unset, not a
+ * number, or {@code 0}. A few seconds, so slow input shows up without
waiting for the buffer to
+ * fill. A negative value disables the interval flush.
+ */
+ static final int DEFAULT_FILE_FLUSH_INTERVAL_MS = 5000;
+
public int getFlushInterval() {
- String maxStreamLife =
variables.getVariable("HOP_FILE_OUTPUT_MAX_STREAM_LIFE");
- int flushInterval = 0;
- if (maxStreamLife != null) {
- try {
- flushInterval = Integer.parseInt(maxStreamLife);
- } catch (Exception ex) {
- // Do nothing
- }
+ String maxStreamLife =
variables.getVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE);
+ int flushInterval = Const.toInt(maxStreamLife,
DEFAULT_FILE_FLUSH_INTERVAL_MS);
+ // 0 is what existing hop-config.json files store and means "not
configured".
+ // Only a negative value disables the interval flush.
+ if (flushInterval == 0) {
+ return DEFAULT_FILE_FLUSH_INTERVAL_MS;
}
return flushInterval;
}
@@ -370,7 +374,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
data.splitnr++;
data.fos = null;
data.out = null;
- data.writer = null;
+ assignWriter(null, null);
filename = getOutputFileName(null);
isWriteHeader = isWriteHeader(filename);
initFileStreamWriter(filename);
@@ -388,16 +392,16 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
int flushInterval = getFlushInterval();
if (flushInterval > 0) {
- long currentTime = new Date().getTime();
+ long currentTime = currentFlushTimeMillis();
if (data.lastFileFlushTime == 0) {
data.lastFileFlushTime = currentTime;
- } else if (data.lastFileFlushTime - currentTime > flushInterval) {
+ } else if (currentTime - data.lastFileFlushTime > flushInterval) {
try {
data.getFileStreamsCollection().flushOpenFiles(false);
} catch (IOException e) {
throw new HopException("Unable to flush open files", e);
}
- data.lastFileFlushTime = new Date().getTime();
+ data.lastFileFlushTime = currentTime;
}
}
return true;
@@ -452,6 +456,11 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
}
}
+ /** Clock for the file-flush interval. Tests advance this instead of
sleeping. */
+ protected long currentFlushTimeMillis() {
+ return System.currentTimeMillis();
+ }
+
public void flushOpenFiles(boolean closeAfterFlush) throws IOException {
TextFileOutputData.IFileStreamsCollection coll =
data.getFileStreamsCollection();
if (coll != null) {
@@ -525,6 +534,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
}
incrementLinesOutput();
+ markCurrentFileDirty();
} catch (Exception e) {
throw new HopTransformException("Error writing line", e);
@@ -732,6 +742,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
if (sLine != null && !sLine.trim().isEmpty()) {
data.writer.write(getBinaryString(sLine));
incrementLinesOutput();
+ markCurrentFileDirty();
}
} catch (Exception e) {
logError("Error writing ended tag line: " + e.toString());
@@ -810,9 +821,28 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
retval = true;
}
incrementLinesOutput();
+ markCurrentFileDirty();
return retval;
}
+ /**
+ * Writer and the stream it belongs to move together so dirty-marking does
not scan open files.
+ */
+ private void assignWriter(OutputStream writer, TextFileOutputData.FileStream
fileStream) {
+ data.writer = writer;
+ data.currentFileStream = fileStream;
+ }
+
+ /**
+ * An interval flush clears the dirty flag. Later rows still land in the
buffer, so the flag has
+ * to be set again or the next flush is skipped.
+ */
+ private void markCurrentFileDirty() {
+ if (data.currentFileStream != null) {
+ data.currentFileStream.setDirty(true);
+ }
+ }
+
public String buildFilename(String filename, boolean ziparchive) {
return meta.buildFilename(
filename,
@@ -964,7 +994,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
coll.closeStream(data.writer);
}
}
- data.writer = null;
+ assignWriter(null, null);
data.out = null;
data.fos = null;
if (isDebug()) {
@@ -974,7 +1004,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
} catch (Exception e) {
logError("Exception trying to close file: " + e.toString());
setErrors(1);
- data.writer = null;
+ assignWriter(null, null);
data.out = null;
data.fos = null;
retval = false;
@@ -1112,7 +1142,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
}
coll.flushOpenFiles(true);
}
- data.writer = null;
+ assignWriter(null, null);
}
private void emitWriteLineageForOpenStream(String filename,
TextFileOutputData.FileStream fs) {
@@ -1150,7 +1180,7 @@ public class TextFileOutput extends
BaseTransform<TextFileOutputMeta, TextFileOu
logError("Unexpected error closing file", e);
setErrors(1);
}
- data.writer = null;
+ assignWriter(null, null);
data.out = null;
data.fos = null;
diff --git
a/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputData.java
b/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputData.java
index 057b84381d..daaf7302a2 100644
---
a/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputData.java
+++
b/plugins/transforms/textfile/src/main/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputData.java
@@ -126,13 +126,18 @@ public class TextFileOutputData extends BaseTransformData
implements ITransformD
}
public void flush() throws IOException {
- if (isDirty) {
+ if (isDirty && getBufferedOutputStream() != null) {
getBufferedOutputStream().flush();
isDirty = false;
}
}
public void close() throws IOException {
+ // Flush even when the dirty flag was cleared. Otherwise bytes written
after the last
+ // interval flush are dropped when the buffer is discarded below.
+ if (getBufferedOutputStream() != null) {
+ getBufferedOutputStream().flush();
+ }
setBufferedOutputStream(null);
getCompressedOutputStream().close();
setCompressedOutputStream(null);
@@ -210,16 +215,17 @@ public class TextFileOutputData extends BaseTransformData
implements ITransformD
@Override
public void flushOpenFiles(boolean closeAfterFlush) throws IOException {
for (FileStream outputStream : streamsList) {
- if (outputStream.isDirty()) {
- try {
+ try {
+ if (outputStream.isDirty()) {
outputStream.flush();
- if (closeAfterFlush && outputStream.isOpen()) {
- outputStream.close();
- numOpenFiles--;
- }
- } catch (IOException e) {
- e.printStackTrace();
}
+ // An interval flush clears the dirty flag. End-of-run still has to
close the stream.
+ if (closeAfterFlush && outputStream.isOpen()) {
+ outputStream.close();
+ numOpenFiles--;
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
}
}
}
@@ -404,15 +410,18 @@ public class TextFileOutputData extends BaseTransformData
implements ITransformD
@Override
public void flushOpenFiles(boolean closeAfterFlush) {
for (FileStreamsCollectionEntry collectionEntry : indexMap.values()) {
- if (collectionEntry.getFileStream().isDirty()) {
- try {
- collectionEntry.getFileStream().flush();
- if (closeAfterFlush) {
- collectionEntry.getFileStream().close();
- }
- } catch (IOException e) {
- e.printStackTrace();
+ FileStream fileStream = collectionEntry.getFileStream();
+ try {
+ if (fileStream.isDirty()) {
+ fileStream.flush();
}
+ // An interval flush clears the dirty flag. End-of-run still has to
close the stream.
+ if (closeAfterFlush && fileStream.isOpen()) {
+ fileStream.close();
+ numOpenFiles--;
+ }
+ } catch (IOException e) {
+ e.printStackTrace();
}
}
}
@@ -455,6 +464,9 @@ public class TextFileOutputData extends BaseTransformData
implements ITransformD
public OutputStream writer;
+ /** Stream {@link #writer} currently writes to, so a row can mark that file
dirty directly. */
+ public FileStream currentFileStream;
+
public DecimalFormat defaultDecimalFormat;
public DecimalFormatSymbols defaultDecimalFormatSymbols;
diff --git
a/plugins/transforms/textfile/src/test/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputFlushTest.java
b/plugins/transforms/textfile/src/test/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputFlushTest.java
new file mode 100644
index 0000000000..6906371643
--- /dev/null
+++
b/plugins/transforms/textfile/src/test/java/org/apache/hop/pipeline/transforms/textfileoutput/TextFileOutputFlushTest.java
@@ -0,0 +1,309 @@
+/*
+ * 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.textfileoutput;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.BufferedOutputStream;
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.io.OutputStream;
+import java.nio.charset.StandardCharsets;
+import java.util.zip.GZIPInputStream;
+import org.apache.hop.core.Const;
+import org.apache.hop.core.compress.CompressionOutputStream;
+import org.apache.hop.core.compress.CompressionPluginType;
+import org.apache.hop.core.compress.gzip.GzipCompressionProvider;
+import org.apache.hop.core.plugins.PluginRegistry;
+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.IVariables;
+import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
+import org.apache.hop.pipeline.transforms.mock.TransformMockHelper;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.mockito.Mockito;
+
+class TextFileOutputFlushTest {
+ @RegisterExtension
+ static RestoreHopEngineEnvironmentExtension env = new
RestoreHopEngineEnvironmentExtension();
+
+ private TransformMockHelper<TextFileOutputMeta, TextFileOutputData> helper;
+
+ @BeforeAll
+ static void setUpBeforeClass() throws Exception {
+ PluginRegistry.addPluginType(CompressionPluginType.getInstance());
+ PluginRegistry.init();
+ }
+
+ @BeforeEach
+ void setUp() {
+ helper =
+ new TransformMockHelper<>(
+ "Text file output flush", TextFileOutputMeta.class,
TextFileOutputData.class);
+ Mockito.when(helper.logChannelFactory.create(Mockito.any(), Mockito.any()))
+ .thenReturn(helper.iLogChannel);
+ Mockito.when(helper.pipeline.isRunning()).thenReturn(true);
+ }
+
+ @AfterEach
+ void tearDown() {
+ helper.cleanUp();
+ }
+
+ @Test
+ void flushIntervalDefaultsToAFewSeconds() throws Exception {
+ FlushProbe transform = newProbe();
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, null);
+ assertEquals(TextFileOutput.DEFAULT_FILE_FLUSH_INTERVAL_MS,
transform.getFlushInterval());
+
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "0");
+ assertEquals(TextFileOutput.DEFAULT_FILE_FLUSH_INTERVAL_MS,
transform.getFlushInterval());
+
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "nope");
+ assertEquals(TextFileOutput.DEFAULT_FILE_FLUSH_INTERVAL_MS,
transform.getFlushInterval());
+
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "2500");
+ assertEquals(2500, transform.getFlushInterval());
+
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "-1");
+ assertEquals(-1, transform.getFlushInterval());
+ }
+
+ @Test
+ void negativeIntervalDoesNotFlushOnTheClock() throws Exception {
+ FlushProbe transform = newProbe();
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "-1");
+ assertTrue(transform.init());
+
+ transform.now = 1_000_000L;
+ transform.row = new Object[] {"a"};
+ assertTrue(transform.processRow());
+
+ transform.now = 1_060_000L;
+ transform.row = new Object[] {"b"};
+ assertTrue(transform.processRow());
+ assertEquals("", transform.written(), "a negative interval does not flush
on the clock");
+ }
+
+ @Test
+ void intervalFlushAfterLastRowStillClosesGzipFile() throws Exception {
+ FlushProbe transform = newProbe("GZip");
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "1000");
+ assertTrue(transform.init());
+
+ transform.now = 1_000_000L;
+ transform.row = new Object[] {"a"};
+ assertTrue(transform.processRow());
+
+ // Last data row. The interval elapses while writing it, so the flush
clears the dirty flag.
+ transform.now = 1_002_000L;
+ transform.row = new Object[] {"b"};
+ assertTrue(transform.processRow());
+ assertFalse(transform.currentStreamDirty(), "interval flush cleared the
dirty flag");
+ assertTrue(transform.currentStreamOpen(), "interval flush must not close
the file");
+ assertFalse(transform.isOutputClosed());
+
+ transform.row = null;
+ assertFalse(transform.processRow());
+ assertFalse(transform.currentStreamOpen());
+ assertTrue(transform.isOutputClosed());
+ assertEquals("a\nb\n", gunzip(transform.gzipBytes()));
+ }
+
+ @Test
+ void closeAfterFlushClosesStreamThatIntervalFlushAlreadyCleaned() throws
Exception {
+ TextFileOutputData data = new TextFileOutputData();
+ assertCleanOpenStreamIsClosed(data.new FileStreamsList());
+ assertCleanOpenStreamIsClosed(data.new FileStreamsMap());
+ }
+
+ private static void
assertCleanOpenStreamIsClosed(TextFileOutputData.IFileStreamsCollection coll)
+ throws Exception {
+ String collection = coll.getClass().getSimpleName();
+ ByteArrayOutputStream raw = new ByteArrayOutputStream();
+ CompressionOutputStream compression = new
GzipCompressionProvider().createOutputStream(raw);
+ BufferedOutputStream buffered = new BufferedOutputStream(compression,
5000);
+ TextFileOutputData.FileStream stream =
+ new TextFileOutputData().new FileStream(raw, compression, buffered);
+ buffered.write("hello\n".getBytes(StandardCharsets.UTF_8));
+ stream.setDirty(true);
+ coll.add("out.txt", stream);
+
+ coll.flushOpenFiles(false);
+ assertFalse(stream.isDirty(), collection);
+ assertTrue(stream.isOpen(), collection);
+ assertEquals(1, coll.getNumOpenFiles(), collection);
+
+ coll.flushOpenFiles(true);
+ assertFalse(stream.isOpen(), collection);
+ assertEquals(0, coll.getNumOpenFiles(), collection);
+ assertEquals("hello\n", gunzip(raw.toByteArray()), collection);
+ }
+
+ private static String gunzip(byte[] gzipBytes) throws IOException {
+ try (GZIPInputStream in = new GZIPInputStream(new
ByteArrayInputStream(gzipBytes))) {
+ return new String(in.readAllBytes(), StandardCharsets.UTF_8);
+ }
+ }
+
+ @Test
+ void slowRowsAreFlushedOnTheIntervalAndAgainAfterThat() throws Exception {
+ FlushProbe transform = newProbe();
+ transform.setVariable(Const.HOP_FILE_OUTPUT_MAX_STREAM_LIFE, "1000");
+ assertTrue(transform.init());
+
+ transform.now = 1_000_000L;
+ transform.row = new Object[] {"a"};
+ assertTrue(transform.processRow());
+ assertEquals("", transform.written(), "first row stays buffered until the
interval elapses");
+
+ transform.now = 1_000_500L;
+ transform.row = new Object[] {"b"};
+ assertTrue(transform.processRow());
+ assertEquals("", transform.written(), "a row inside the interval does not
flush");
+
+ transform.now = 1_001_001L;
+ transform.row = new Object[] {"c"};
+ assertTrue(transform.processRow());
+ assertEquals("a\nb\nc\n", transform.written());
+
+ transform.now = 1_002_002L;
+ transform.row = new Object[] {"d"};
+ assertTrue(transform.processRow());
+ assertEquals("a\nb\nc\nd\n", transform.written());
+ }
+
+ private FlushProbe newProbe() {
+ return newProbe("None");
+ }
+
+ private FlushProbe newProbe(String compression) {
+ TextFileOutputMeta meta = new TextFileOutputMeta();
+ meta.setDefault();
+ meta.setFileCompression(compression);
+ meta.setHeaderEnabled(false);
+ meta.setFooterEnabled(false);
+ meta.setSeparator("");
+ meta.setEnclosure("");
+ meta.setFileFormat("UNIX");
+ meta.setEncoding(Const.UTF_8);
+ meta.setCreateParentFolder(false);
+ meta.setEndedLine(null);
+ meta.getFileSettings().setFileName("out.txt");
+ meta.getFileSettings().setExtension("");
+ meta.getFileSettings().setDoNotOpenNewFileInit(true);
+ meta.getFileSettings().setAddToResultFiles(false);
+ meta.getFileSettings().setFastDump(true);
+
+ FlushProbe transform =
+ new FlushProbe(
+ helper.transformMeta,
+ meta,
+ new TextFileOutputData(),
+ helper.pipelineMeta,
+ helper.pipeline);
+ RowMeta rowMeta = new RowMeta();
+ ValueMetaString column = new ValueMetaString("col");
+ column.setLength(-1);
+ rowMeta.addValueMeta(column);
+ transform.setInputRowMeta(rowMeta);
+ return transform;
+ }
+
+ private static final class FlushProbe extends TextFileOutput {
+ private long now;
+ private Object[] row;
+ private boolean outputClosed;
+ private final ByteArrayOutputStream written = new ByteArrayOutputStream();
+
+ private FlushProbe(
+ org.apache.hop.pipeline.transform.TransformMeta transformMeta,
+ TextFileOutputMeta meta,
+ TextFileOutputData data,
+ org.apache.hop.pipeline.PipelineMeta pipelineMeta,
+ org.apache.hop.pipeline.Pipeline pipeline) {
+ super(transformMeta, meta, data, 0, pipelineMeta, pipeline);
+ }
+
+ private String written() {
+ return written.toString(StandardCharsets.UTF_8);
+ }
+
+ private byte[] gzipBytes() {
+ return written.toByteArray();
+ }
+
+ private boolean isOutputClosed() {
+ return outputClosed;
+ }
+
+ private boolean currentStreamOpen() {
+ TextFileOutputData.FileStream last =
data.getFileStreamsCollection().getLastStream();
+ return last != null && last.isOpen();
+ }
+
+ private boolean currentStreamDirty() {
+ TextFileOutputData.FileStream last =
data.getFileStreamsCollection().getLastStream();
+ return last != null && last.isDirty();
+ }
+
+ @Override
+ protected long currentFlushTimeMillis() {
+ return now;
+ }
+
+ @Override
+ public Object[] getRow() {
+ return row;
+ }
+
+ @Override
+ public void putRow(IRowMeta rowMeta, Object[] row) {
+ // The flush test does not chain rows to a downstream transform.
+ }
+
+ @Override
+ protected OutputStream getOutputStream(
+ String vfsFilename, IVariables variables, boolean append) {
+ return new OutputStream() {
+ @Override
+ public void write(int b) {
+ written.write(b);
+ }
+
+ @Override
+ public void write(byte[] b, int off, int len) {
+ written.write(b, off, len);
+ }
+
+ @Override
+ public void close() {
+ outputClosed = true;
+ }
+ };
+ }
+ }
+}