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