This is an automated email from the ASF dual-hosted git repository.

hansva 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 2873d19146 fix snowflake bulkloader, fixes #2298 (#8089)
2873d19146 is described below

commit 2873d19146a40a9aada5f7b824e4561eecb64ab1
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Tue Aug 25 15:23:03 2026 +0200

    fix snowflake bulkloader, fixes #2298 (#8089)
---
 .../snowflake/bulkloader/SnowflakeBulkLoader.java  |  82 +++++++++++++++-
 .../bulkloader/SnowflakeBulkLoaderData.java        |   5 +
 .../bulkloader/SnowflakeBulkLoaderMeta.java        |  11 +++
 .../bulkloader/SnowflakeBulkLoaderTest.java        | 105 +++++++++++++++++++++
 4 files changed, 198 insertions(+), 5 deletions(-)

diff --git 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoader.java
 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoader.java
index 5c68bab04f..297ab5394e 100644
--- 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoader.java
+++ 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoader.java
@@ -158,6 +158,25 @@ public class SnowflakeBulkLoader
               CONST_FIELD + meta.getJsonField() + "] couldn't be found in the 
input stream!");
         }
         data.fieldnrs.put("json", streamFieldLocation);
+      } else {
+        // No fields were specified: every field on the stream is written, in 
stream order.  The
+        // COPY statement maps those to the table columns by position, so 
describe the table to
+        // find out in which column every value lands.
+        //
+        try {
+          getDbFields();
+        } catch (HopException e) {
+          // Without the table layout we can still write the file, dates 
simply fall back to the
+          // timestamp format.
+          //
+          logBasic("Unable to get the fields of the target table, using 
default date formats", e);
+          data.dbFields = null;
+        }
+      }
+
+      if (meta.getDataTypeId() == SnowflakeBulkLoaderMeta.DATA_TYPE_CSV
+          && !meta.isSpecifyFields()) {
+        data.writeValueMetas = getWriteValueMetas(data.outputRowMeta);
       }
     }
 
@@ -346,6 +365,59 @@ public class SnowflakeBulkLoader
     data.db.execStatement("commit");
   }
 
+  /**
+   * Determines the value metadata to use when writing the fields of the 
stream to a temp file. Date
+   * and timestamp values are written with the conversion mask matching the 
file format of the COPY
+   * statement instead of the mask defined on the stream, since Snowflake only 
parses the format it
+   * was told to expect.
+   *
+   * @param rowMeta The metadata of the rows on the stream
+   * @return The value metadata to write every field of the row with
+   */
+  IValueMeta[] getWriteValueMetas(IRowMeta rowMeta) {
+    IValueMeta[] writeValueMetas = new IValueMeta[rowMeta.size()];
+    for (int i = 0; i < rowMeta.size(); i++) {
+      writeValueMetas[i] = getWriteValueMeta(rowMeta.getValueMeta(i), i);
+    }
+    return writeValueMetas;
+  }
+
+  /**
+   * Determines the value metadata to use for a single field of the stream.
+   *
+   * @param valueMeta The metadata of the field on the stream
+   * @param index The position of the field, which is the position of the 
column it is loaded into
+   * @return The original value metadata, or a copy of it carrying the 
Snowflake conversion mask
+   */
+  private IValueMeta getWriteValueMeta(IValueMeta valueMeta, int index) {
+    if (!valueMeta.isDate()) {
+      return valueMeta;
+    }
+
+    // Without the table layout we don't know the type of the target column, 
the timestamp format
+    // is the safest default as it also carries the date.
+    //
+    String mask = SnowflakeBulkLoaderMeta.TIMESTAMP_MASK;
+    if (data.dbFields != null && index < data.dbFields.size()) {
+      String type = data.dbFields.get(index)[1].toUpperCase();
+      if (type.startsWith("TIMESTAMP")) {
+        mask = SnowflakeBulkLoaderMeta.TIMESTAMP_MASK;
+      } else if (type.startsWith("DATE")) {
+        mask = SnowflakeBulkLoaderMeta.DATE_MASK;
+      } else if (type.startsWith("TIME")) {
+        mask = SnowflakeBulkLoaderMeta.TIME_MASK;
+      } else {
+        // The value doesn't end up in a date column, leave the format of the 
stream alone.
+        //
+        return valueMeta;
+      }
+    }
+
+    IValueMeta writeValueMeta = valueMeta.clone();
+    writeValueMeta.setConversionMask(mask);
+    return writeValueMeta;
+  }
+
   /**
    * Writes an individual row of data to a temp file
    *
@@ -353,7 +425,7 @@ public class SnowflakeBulkLoader
    * @param row The input row
    * @throws HopTransformException
    */
-  private void writeRowToFile(IRowMeta rowMeta, Object[] row) throws 
HopTransformException {
+  void writeRowToFile(IRowMeta rowMeta, Object[] row) throws 
HopTransformException {
     try {
       if (meta.getDataTypeId() == SnowflakeBulkLoaderMeta.DATA_TYPE_CSV
           && !meta.isSpecifyFields()) {
@@ -364,7 +436,7 @@ public class SnowflakeBulkLoader
           if (i > 0 && data.binarySeparator.length > 0) {
             data.writer.write(data.binarySeparator);
           }
-          IValueMeta v = rowMeta.getValueMeta(i);
+          IValueMeta v = data.writeValueMetas[i];
           Object valueData = row[i];
 
           // no special null value default was specified since no fields are 
specified at all
@@ -388,13 +460,13 @@ public class SnowflakeBulkLoader
 
             if (field[1].toUpperCase().startsWith("TIMESTAMP")) {
               v = new ValueMetaDate();
-              v.setConversionMask("yyyy-MM-dd HH:mm:ss.SSS");
+              v.setConversionMask(SnowflakeBulkLoaderMeta.TIMESTAMP_MASK);
             } else if (field[1].toUpperCase().startsWith("DATE")) {
               v = new ValueMetaDate();
-              v.setConversionMask("yyyy-MM-dd");
+              v.setConversionMask(SnowflakeBulkLoaderMeta.DATE_MASK);
             } else if (field[1].toUpperCase().startsWith("TIME")) {
               v = new ValueMetaDate();
-              v.setConversionMask("HH:mm:ss.SSS");
+              v.setConversionMask(SnowflakeBulkLoaderMeta.TIME_MASK);
             } else if (field[1].toUpperCase().startsWith("NUMBER")
                 || field[1].toUpperCase().startsWith("FLOAT")) {
               v = new ValueMetaBigNumber();
diff --git 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderData.java
 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderData.java
index 95880c9d9f..6bfbf42a84 100644
--- 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderData.java
+++ 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderData.java
@@ -25,6 +25,7 @@ import org.apache.hop.core.compress.CompressionOutputStream;
 import org.apache.hop.core.database.Database;
 import org.apache.hop.core.database.DatabaseMeta;
 import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.IValueMeta;
 import org.apache.hop.pipeline.transform.BaseTransformData;
 import org.apache.hop.pipeline.transform.ITransformData;
 
@@ -62,6 +63,10 @@ public class SnowflakeBulkLoaderData extends 
BaseTransformData implements ITrans
   // The metadata about the output row
   public IRowMeta outputRowMeta;
 
+  // The metadata used to write every field of the output row to the temp 
files.  Dates and
+  // timestamps carry the conversion mask Snowflake expects instead of the one 
from the stream.
+  public IValueMeta[] writeValueMetas;
+
   // Byte arrays for constant characters put into output files.
   public byte[] binarySeparator;
   public byte[] binaryEnclosure;
diff --git 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderMeta.java
 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderMeta.java
index 5af71d0c5b..0f39451413 100644
--- 
a/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderMeta.java
+++ 
b/plugins/databases/snowflake/src/main/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderMeta.java
@@ -78,6 +78,16 @@ public class SnowflakeBulkLoaderMeta
   public static final String ENCLOSURE = "\"";
   public static final String DATE_FORMAT_STRING = "yyyy-MM-dd";
   public static final String TIMESTAMP_FORMAT_STRING = "YYYY-MM-DD 
HH24:MI:SS.FF3";
+  public static final String TIME_FORMAT_STRING = "HH24:MI:SS.FF3";
+
+  /*
+   * The Java conversion masks matching the Snowflake file formats above.  
Dates and timestamps have
+   * to be written to the temp files in exactly the format the COPY statement 
declares, otherwise
+   * Snowflake refuses to parse them.
+   */
+  public static final String DATE_MASK = "yyyy-MM-dd";
+  public static final String TIMESTAMP_MASK = "yyyy-MM-dd HH:mm:ss.SSS";
+  public static final String TIME_MASK = "HH:mm:ss.SSS";
 
   /** The valid location type codes */
   public static final String[] LOCATION_TYPE_CODES = {"user", "table", 
"internal_stage"};
@@ -1058,6 +1068,7 @@ public class SnowflakeBulkLoaderMeta
       returnValue.append("ESCAPE_UNENCLOSED_FIELD = '\\\\' 
FIELD_OPTIONALLY_ENCLOSED_BY='\"' ");
       returnValue.append("SKIP_HEADER = 0 DATE_FORMAT = 
'").append(DATE_FORMAT_STRING).append("' ");
       returnValue.append("TIMESTAMP_FORMAT = 
'").append(TIMESTAMP_FORMAT_STRING).append("' ");
+      returnValue.append("TIME_FORMAT = 
'").append(TIME_FORMAT_STRING).append("' ");
       returnValue.append("TRIM_SPACE = ").append(trimWhitespace).append(" ");
       if (!StringUtils.isEmpty(nullIf)) {
         returnValue.append("NULL_IF = (");
diff --git 
a/plugins/databases/snowflake/src/test/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderTest.java
 
b/plugins/databases/snowflake/src/test/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderTest.java
index d4d8f13e37..cfc5e7dd5a 100644
--- 
a/plugins/databases/snowflake/src/test/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderTest.java
+++ 
b/plugins/databases/snowflake/src/test/java/org/apache/hop/pipeline/transforms/snowflake/bulkloader/SnowflakeBulkLoaderTest.java
@@ -17,7 +17,9 @@
 
 package org.apache.hop.pipeline.transforms.snowflake.bulkloader;
 
+import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
@@ -28,9 +30,21 @@ import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.io.ByteArrayOutputStream;
+import java.nio.charset.StandardCharsets;
 import java.sql.Connection;
+import java.util.ArrayList;
+import java.util.Calendar;
+import java.util.Date;
+import java.util.GregorianCalendar;
 import org.apache.hop.core.database.Database;
 import org.apache.hop.core.database.DatabaseMeta;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.IValueMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaDate;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.core.row.value.ValueMetaTimestamp;
 import org.apache.hop.pipeline.PipelineMeta;
 import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
 import org.apache.hop.pipeline.transform.TransformMeta;
@@ -188,4 +202,95 @@ class SnowflakeBulkLoaderTest {
 
     verify(bulkLoaderSpy, never()).truncateTable();
   }
+
+  @Test
+  void testDateIsWrittenInTheSnowflakeTimestampFormatWhenTheTableIsUnknown() {
+    data.dbFields = null;
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(new ValueMetaDate("created"));
+
+    IValueMeta[] writeValueMetas = bulkLoaderSpy.getWriteValueMetas(rowMeta);
+
+    assertEquals(SnowflakeBulkLoaderMeta.TIMESTAMP_MASK, 
writeValueMetas[0].getConversionMask());
+  }
+
+  @Test
+  void testDateIsWrittenInTheFormatOfTheTargetColumn() {
+    data.dbFields = new ArrayList<>();
+    data.dbFields.add(new String[] {"BIRTHDAY", "DATE"});
+    data.dbFields.add(new String[] {"CREATED", "TIMESTAMP_NTZ(9)"});
+    data.dbFields.add(new String[] {"STARTED", "TIME(9)"});
+
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(new ValueMetaDate("birthday"));
+    rowMeta.addValueMeta(new ValueMetaTimestamp("created"));
+    rowMeta.addValueMeta(new ValueMetaDate("started"));
+
+    IValueMeta[] writeValueMetas = bulkLoaderSpy.getWriteValueMetas(rowMeta);
+
+    assertEquals(SnowflakeBulkLoaderMeta.DATE_MASK, 
writeValueMetas[0].getConversionMask());
+    assertEquals(SnowflakeBulkLoaderMeta.TIMESTAMP_MASK, 
writeValueMetas[1].getConversionMask());
+    assertEquals(SnowflakeBulkLoaderMeta.TIME_MASK, 
writeValueMetas[2].getConversionMask());
+  }
+
+  @Test
+  void testFormatOfTheStreamIsKeptForNonDateColumns() {
+    data.dbFields = new ArrayList<>();
+    data.dbFields.add(new String[] {"CREATED", "VARCHAR(16777216)"});
+    data.dbFields.add(new String[] {"NAME", "VARCHAR(16777216)"});
+
+    IRowMeta rowMeta = new RowMeta();
+    IValueMeta date = new ValueMetaDate("created");
+    date.setConversionMask("dd/MM/yyyy");
+    rowMeta.addValueMeta(date);
+    IValueMeta name = new ValueMetaString("name");
+    rowMeta.addValueMeta(name);
+
+    IValueMeta[] writeValueMetas = bulkLoaderSpy.getWriteValueMetas(rowMeta);
+
+    assertSame(date, writeValueMetas[0]);
+    assertSame(name, writeValueMetas[1]);
+  }
+
+  @Test
+  void testTheFormatOfTheStreamIsNotModified() {
+    data.dbFields = null;
+    IRowMeta rowMeta = new RowMeta();
+    IValueMeta date = new ValueMetaDate("created");
+    date.setConversionMask("dd/MM/yyyy");
+    rowMeta.addValueMeta(date);
+
+    bulkLoaderSpy.getWriteValueMetas(rowMeta);
+
+    assertEquals("dd/MM/yyyy", date.getConversionMask());
+  }
+
+  @Test
+  void testDatesAreWrittenToTheTempFileInTheSnowflakeFormat() throws Exception 
{
+    data.dbFields = new ArrayList<>();
+    data.dbFields.add(new String[] {"BIRTHDAY", "DATE"});
+    data.dbFields.add(new String[] {"CREATED", "TIMESTAMP_NTZ(9)"});
+
+    IRowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(new ValueMetaDate("birthday"));
+    rowMeta.addValueMeta(new ValueMetaDate("created"));
+
+    data.outputRowMeta = rowMeta;
+    data.writeValueMetas = bulkLoaderSpy.getWriteValueMetas(rowMeta);
+    data.binarySeparator = 
SnowflakeBulkLoaderMeta.CSV_DELIMITER.getBytes(StandardCharsets.UTF_8);
+    data.binaryEnclosure = 
SnowflakeBulkLoaderMeta.ENCLOSURE.getBytes(StandardCharsets.UTF_8);
+    data.escapeCharacters =
+        
SnowflakeBulkLoaderMeta.CSV_ESCAPE_CHAR.getBytes(StandardCharsets.UTF_8);
+    data.binaryNewline =
+        
SnowflakeBulkLoaderMeta.CSV_RECORD_DELIMITER.getBytes(StandardCharsets.UTF_8);
+    ByteArrayOutputStream output = new ByteArrayOutputStream();
+    data.writer = output;
+
+    Date date = new GregorianCalendar(2023, Calendar.FEBRUARY, 9, 10, 11, 
12).getTime();
+    bulkLoaderSpy.writeRowToFile(rowMeta, new Object[] {date, date});
+
+    assertEquals(
+        "2023-02-09,2023-02-09 10:11:12.000" + 
SnowflakeBulkLoaderMeta.CSV_RECORD_DELIMITER,
+        output.toString(StandardCharsets.UTF_8));
+  }
 }

Reply via email to