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;
+        }
+      };
+    }
+  }
+}

Reply via email to