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 01708bd4c8 Issue #3864 : Select logical types on Parquet output fields
(#8684)
01708bd4c8 is described below
commit 01708bd4c8e8887278804f5fe7bab627cf54daa6
Author: Matt Casters <[email protected]>
AuthorDate: Thu Oct 1 20:59:07 2026 +0200
Issue #3864 : Select logical types on Parquet output fields (#8684)
* Issue #3864 : Select logical types on Parquet output fields
* Issue #3864 : Propose TimestampMillis for a Hop Date in Parquet Get Fields
---
.../pipeline/transforms/parquet-file-output.adoc | 101 ++++-
.../parquet/transforms/output/ParquetField.java | 73 ++--
.../transforms/output/ParquetFieldType.java | 345 +++++++++++++++
.../parquet/transforms/output/ParquetOutput.java | 36 +-
.../transforms/output/ParquetOutputDialog.java | 70 ++-
.../transforms/output/ParquetWriteSupport.java | 79 ++--
.../output/messages/messages_en_US.properties | 6 +
.../transforms/output/ParquetFieldTest.java | 34 ++
.../transforms/output/ParquetLogicalTypeTest.java | 475 +++++++++++++++++++++
.../transforms/output/ParquetOutputMetaTest.java | 22 +-
.../transforms/output/ParquetOutputTest.java | 11 +-
11 files changed, 1181 insertions(+), 71 deletions(-)
diff --git
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-output.adoc
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-output.adoc
index 1b937e0b9c..6f2208c156 100644
---
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-output.adoc
+++
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-output.adoc
@@ -174,13 +174,110 @@
image:tech/parquet/parquet-output-dialog-fields-tab.png[Parquet file output Fiel
|Fields
|You can specify which fields to write and in which order.
You can use the "Get Fields" button to populate the dialog.
-Leave empty to output all input fields.
+Leave the list empty to output all input fields with the types in
<<hop-to-parquet-types,Hop to Parquet types>>.
+
+|Source field
+|Incoming field to write.
+
+|Target field
+|Column name in the Parquet file.
+Leave empty to use the source field name.
+
+|Parquet type
+|https://parquet.apache.org/docs/file-format/types/logicaltypes/[Logical or
physical type] of the column.
+Leave empty to keep the type chosen from the Hop type.
+Get Fields proposes a type for each field: a Date is proposed as
`TimestampMillis`, the same TIMESTAMP(MILLIS) written when the Parquet type is
left empty, and a Timestamp is proposed as `TimestampMicros`. Select `Date` for
a calendar date, for example so BigQuery loads the column as a DATE.
+Change the proposed type when you need a different one.
+
+|Precision
+|Total number of digits.
+Used only when the Parquet type is `Decimal`.
+Leave empty to use the length of the source field.
+
+|Scale
+|Digits after the decimal point.
+Used only when the Parquet type is `Decimal`.
+Leave empty to use the precision of the source field.
+|===
+
+=== Choosing a Parquet type
+[options="header"]
|===
+|Parquet type|Column|Value
+
+|UTF8
+|STRING
+|Text of the value
+
+|Boolean
+|BOOLEAN
+|Boolean value
+
+|Int32
+|INT32
+|Whole number. A value outside the 32-bit range stops the transform.
+
+|Int64
+|INT64
+|Whole number
+
+|Float
+|FLOAT
+|32-bit floating point number
+
+|Double
+|DOUBLE
+|64-bit floating point number
+
+|Binary
+|BINARY
+|Bytes of the value
+
+|Date
+|DATE
+|Calendar date in the JVM time zone. The time of day is dropped.
+BigQuery and other warehouses load this logical type as a DATE.
+Parquet Input reads it back as midnight of that date in the same time zone.
+
+|TimeMillis
+|TIME(MILLIS), not adjusted to UTC
+|Time of day in the JVM time zone, with millisecond precision
+
+|TimeMicros
+|TIME(MICROS), not adjusted to UTC
+|Time of day in the JVM time zone, with microsecond precision
+
+|TimestampMillis
+|TIMESTAMP(MILLIS), adjusted to UTC
+|Instant, with millisecond precision
+
+|TimestampMicros
+|TIMESTAMP(MICROS), adjusted to UTC
+|Instant, with microsecond precision. The nanoseconds beyond the microsecond
are dropped.
+
+|Decimal
+|DECIMAL(precision, scale)
+|Exact decimal, rounded with the rounding type of the field.
+Precision is 1 to 38 and scale is 0 to the precision, taken from the Precision
and Scale columns or, when those are empty, from the length and precision of
the source field.
+A value with more digits before the decimal point than the column holds stops
the transform.
+
+|JSON
+|JSON
+|JSON text
+
+|UUID
+|UUID
+|The 16 bytes of the UUID. The value is the text form, for example
`00112233-4455-6677-8899-aabbccddeeff`.
+|===
+
+A field listed without a Parquet type, and every field written because the
list is empty, uses the mapping below.
+An existing pipeline that does not set a Parquet type keeps that mapping: a
Date stays a TIMESTAMP(MILLIS) until you select `Date`.
+[#hop-to-parquet-types]
== Hop to Parquet types
-Every field is written as an optional column of the
https://parquet.apache.org/docs/file-format/types/logicaltypes/[Parquet type]
below.
+When the Parquet type is left empty, every field is written as an optional
column of the
https://parquet.apache.org/docs/file-format/types/logicaltypes/[Parquet type]
below.
[options="header"]
|===
diff --git
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetField.java
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetField.java
index 5ff0b477c5..2c8f0fa85d 100644
---
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetField.java
+++
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetField.java
@@ -17,8 +17,15 @@
package org.apache.hop.parquet.transforms.output;
+import java.util.Arrays;
+import lombok.Getter;
+import lombok.Setter;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.util.Utils;
import org.apache.hop.metadata.api.HopMetadataProperty;
+@Getter
+@Setter
public class ParquetField {
@HopMetadataProperty(key = "source_field")
private String sourceFieldName;
@@ -26,6 +33,21 @@ public class ParquetField {
@HopMetadataProperty(key = "target_field")
private String targetFieldName;
+ /**
+ * Parquet type code ({@link ParquetFieldType#getCode()}). Empty keeps the
type chosen from the
+ * Hop type of the source field.
+ */
+ @HopMetadataProperty(key = "parquet_type")
+ private String parquetType;
+
+ /** Decimal precision. Used only when the Parquet type is Decimal. */
+ @HopMetadataProperty(key = "precision")
+ private String precision;
+
+ /** Decimal scale. Used only when the Parquet type is Decimal. */
+ @HopMetadataProperty(key = "scale")
+ private String scale;
+
public ParquetField() {}
public ParquetField(String sourceFieldName, String targetFieldName) {
@@ -34,38 +56,31 @@ public class ParquetField {
}
public ParquetField(ParquetField f) {
- this(f.sourceFieldName, f.targetFieldName);
+ this.sourceFieldName = f.sourceFieldName;
+ this.targetFieldName = f.targetFieldName;
+ this.parquetType = f.parquetType;
+ this.precision = f.precision;
+ this.scale = f.scale;
}
/**
- * Gets sourceFieldName
- *
- * @return value of sourceFieldName
+ * @return the selected Parquet type, or null when none is selected
+ * @throws HopException when a type is set but is not one of the supported
types
*/
- public String getSourceFieldName() {
- return sourceFieldName;
- }
-
- /**
- * @param sourceFieldName The sourceFieldName to set
- */
- public void setSourceFieldName(String sourceFieldName) {
- this.sourceFieldName = sourceFieldName;
- }
-
- /**
- * Gets targetFieldName
- *
- * @return value of targetFieldName
- */
- public String getTargetFieldName() {
- return targetFieldName;
- }
-
- /**
- * @param targetFieldName The targetFieldName to set
- */
- public void setTargetFieldName(String targetFieldName) {
- this.targetFieldName = targetFieldName;
+ public ParquetFieldType parquetFieldType() throws HopException {
+ if (Utils.isEmpty(parquetType) || parquetType.trim().isEmpty()) {
+ return null;
+ }
+ ParquetFieldType type = ParquetFieldType.fromCode(parquetType);
+ if (type == null) {
+ throw new HopException(
+ "Parquet type '"
+ + parquetType
+ + "' of field '"
+ + (sourceFieldName == null ? "" : sourceFieldName)
+ + "' is not supported. Supported types: "
+ + String.join(", ", Arrays.asList(ParquetFieldType.codes())));
+ }
+ return type;
}
}
diff --git
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetFieldType.java
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetFieldType.java
new file mode 100644
index 0000000000..4cc6cf2e9c
--- /dev/null
+++
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetFieldType.java
@@ -0,0 +1,345 @@
+/*
+ * 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.parquet.transforms.output;
+
+import java.time.Instant;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.temporal.ChronoField;
+import java.util.Date;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.row.IValueMeta;
+import org.apache.hop.core.row.value.ValueMetaTimestamp;
+import org.apache.hop.core.util.Utils;
+import org.apache.hop.metadata.api.IEnumHasCode;
+import org.apache.parquet.io.api.RecordConsumer;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import
org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
+import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Type;
+import org.apache.parquet.schema.Types;
+
+/**
+ * Parquet type a field can be written as. The code is stored in the pipeline
and shown in the
+ * fields grid.
+ *
+ * <p>Date and the time-of-day types use the calendar fields in the JVM time
zone, the same zone
+ * Parquet Input uses when it reads those columns back. Timestamps are
instants adjusted to UTC.
+ */
+public enum ParquetFieldType implements IEnumHasCode {
+ Utf8("UTF8"),
+ Boolean("Boolean"),
+ Int32("Int32"),
+ Int64("Int64"),
+ Float("Float"),
+ Double("Double"),
+ Binary("Binary"),
+ Date("Date"),
+ TimeMillis("TimeMillis"),
+ TimeMicros("TimeMicros"),
+ TimestampMillis("TimestampMillis"),
+ TimestampMicros("TimestampMicros"),
+ Decimal("Decimal"),
+ Json("JSON"),
+ Uuid("UUID");
+
+ private static final LogicalTypeAnnotation TIMESTAMP_MICROS =
+ LogicalTypeAnnotation.timestampType(true, TimeUnit.MICROS);
+
+ private static final int UUID_LENGTH = 16;
+
+ private final String code;
+
+ ParquetFieldType(String code) {
+ this.code = code;
+ }
+
+ @Override
+ public String getCode() {
+ return code;
+ }
+
+ public static String[] codes() {
+ return IEnumHasCode.getCodes(ParquetFieldType.class);
+ }
+
+ /**
+ * @return the type for this code, or null when the code is empty or not one
of the types
+ */
+ public static ParquetFieldType fromCode(String code) {
+ if (code == null) {
+ return null;
+ }
+ String trimmed = code.trim();
+ if (trimmed.isEmpty()) {
+ return null;
+ }
+ return IEnumHasCode.lookupCode(ParquetFieldType.class, trimmed, null);
+ }
+
+ /**
+ * The type Get Fields proposes. A Hop date is proposed as {@link
#TimestampMillis}, the
+ * TIMESTAMP(MILLIS) the writer uses when the Parquet type is left empty, so
the time of day is
+ * kept. Select {@link #Date} for a calendar date. Every other Hop type
keeps the column the
+ * writer builds when the Parquet type is left empty.
+ *
+ * @param valueMeta the incoming field, or null
+ * @return the proposed type, or null when this Hop type has no Parquet type
+ */
+ public static ParquetFieldType forValueMeta(IValueMeta valueMeta) {
+ if (valueMeta == null) {
+ return null;
+ }
+ return switch (valueMeta.getType()) {
+ case IValueMeta.TYPE_STRING -> Utf8;
+ case IValueMeta.TYPE_BOOLEAN -> Boolean;
+ case IValueMeta.TYPE_INTEGER -> Int64;
+ case IValueMeta.TYPE_NUMBER -> Double;
+ case IValueMeta.TYPE_BIGNUMBER -> decimalFits(valueMeta) ? Decimal :
Utf8;
+ case IValueMeta.TYPE_BINARY -> Binary;
+ case IValueMeta.TYPE_DATE -> TimestampMillis;
+ case IValueMeta.TYPE_TIMESTAMP -> TimestampMicros;
+ case IValueMeta.TYPE_JSON -> Json;
+ case IValueMeta.TYPE_UUID -> Uuid;
+ default -> null;
+ };
+ }
+
+ /**
+ * Same rule as {@link ParquetOutput#avroType}: a big number with a usable
length is a DECIMAL,
+ * otherwise it is written as text.
+ */
+ private static boolean decimalFits(IValueMeta valueMeta) {
+ int length = valueMeta.getLength();
+ int scale = Math.max(valueMeta.getPrecision(), 0);
+ return length > 0 && length <= ParquetOutput.MAX_DECIMAL_PRECISION &&
scale <= length;
+ }
+
+ /** The optional column of this type. */
+ Type column(String name, ParquetField field, IValueMeta valueMeta) throws
HopException {
+ return switch (this) {
+ case Utf8 ->
+ Types.optional(PrimitiveTypeName.BINARY)
+ .as(LogicalTypeAnnotation.stringType())
+ .named(name);
+ case Boolean -> Types.optional(PrimitiveTypeName.BOOLEAN).named(name);
+ case Int32 -> Types.optional(PrimitiveTypeName.INT32).named(name);
+ case Int64 -> Types.optional(PrimitiveTypeName.INT64).named(name);
+ case Float -> Types.optional(PrimitiveTypeName.FLOAT).named(name);
+ case Double -> Types.optional(PrimitiveTypeName.DOUBLE).named(name);
+ case Binary -> Types.optional(PrimitiveTypeName.BINARY).named(name);
+ case Date ->
+
Types.optional(PrimitiveTypeName.INT32).as(LogicalTypeAnnotation.dateType()).named(name);
+ case TimeMillis ->
+ Types.optional(PrimitiveTypeName.INT32)
+ .as(LogicalTypeAnnotation.timeType(false, TimeUnit.MILLIS))
+ .named(name);
+ case TimeMicros ->
+ Types.optional(PrimitiveTypeName.INT64)
+ .as(LogicalTypeAnnotation.timeType(false, TimeUnit.MICROS))
+ .named(name);
+ case TimestampMillis ->
+ Types.optional(PrimitiveTypeName.INT64)
+ .as(LogicalTypeAnnotation.timestampType(true, TimeUnit.MILLIS))
+ .named(name);
+ case TimestampMicros ->
+
Types.optional(PrimitiveTypeName.INT64).as(TIMESTAMP_MICROS).named(name);
+ case Decimal ->
+ Types.optional(PrimitiveTypeName.BINARY).as(decimal(field,
valueMeta)).named(name);
+ case Json ->
+
Types.optional(PrimitiveTypeName.BINARY).as(LogicalTypeAnnotation.jsonType()).named(name);
+ case Uuid ->
+ Types.optional(PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY)
+ .length(UUID_LENGTH)
+ .as(LogicalTypeAnnotation.uuidType())
+ .named(name);
+ };
+ }
+
+ /** Writes one non-null value of this type. The caller has already opened
the field. */
+ void write(RecordConsumer consumer, ParquetField field, IValueMeta
valueMeta, Object valueData)
+ throws HopException {
+ String name = nameOf(field, valueMeta);
+ switch (this) {
+ case Utf8, Json ->
+ consumer.addBinary(
+
org.apache.parquet.io.api.Binary.fromString(requireString(valueMeta,
valueData)));
+ case Boolean -> consumer.addBoolean(requireBoolean(valueMeta,
valueData));
+ case Int32 -> consumer.addInteger(int32(name, valueMeta, valueData));
+ case Int64 -> consumer.addLong(requireLong(valueMeta, valueData));
+ case Float -> consumer.addFloat((float) requireDouble(valueMeta,
valueData));
+ case Double -> consumer.addDouble(requireDouble(valueMeta, valueData));
+ case Binary ->
+ consumer.addBinary(
+ org.apache.parquet.io.api.Binary.fromConstantByteArray(
+ requireBytes(valueMeta, valueData)));
+ case Date -> consumer.addInteger(dateDays(name, valueMeta, valueData));
+ case TimeMillis -> consumer.addInteger(timeOfDayMillis(valueMeta,
valueData));
+ case TimeMicros -> consumer.addLong(timeOfDayMicros(valueMeta,
valueData));
+ case TimestampMillis ->
+ consumer.addLong(ParquetWriteSupport.epochValue(valueMeta,
valueData, null));
+ case TimestampMicros ->
+ consumer.addLong(ParquetWriteSupport.epochValue(valueMeta,
valueData, TIMESTAMP_MICROS));
+ case Decimal ->
+ consumer.addBinary(
+ ParquetWriteSupport.decimalBytes(
+ name, valueMeta, valueMeta.getBigNumber(valueData),
decimal(field, valueMeta)));
+ case Uuid ->
+
consumer.addBinary(ParquetWriteSupport.uuidBytes(requireString(valueMeta,
valueData)));
+ }
+ }
+
+ /**
+ * DECIMAL annotation for this field. An explicit precision and scale win;
otherwise the length
+ * and precision of the source field are used.
+ */
+ DecimalLogicalTypeAnnotation decimal(ParquetField field, IValueMeta
valueMeta)
+ throws HopException {
+ String name = nameOf(field, valueMeta);
+ Integer precisionText =
+ wholeNumber(field == null ? null : field.getPrecision(), name,
"Precision");
+ Integer scaleText = wholeNumber(field == null ? null : field.getScale(),
name, "Scale");
+ int precision = precisionText != null ? precisionText :
valueMeta.getLength();
+ int scale = scaleText != null ? scaleText :
Math.max(valueMeta.getPrecision(), 0);
+ if (precision < 1
+ || precision > ParquetOutput.MAX_DECIMAL_PRECISION
+ || scale < 0
+ || scale > precision) {
+ throw new HopException(
+ "Field '"
+ + name
+ + "' is written as Decimal with precision "
+ + precision
+ + " and scale "
+ + scale
+ + ". Set a precision from 1 to "
+ + ParquetOutput.MAX_DECIMAL_PRECISION
+ + " and a scale from 0 to that precision, or set the length and
precision of the source field.");
+ }
+ return (DecimalLogicalTypeAnnotation)
LogicalTypeAnnotation.decimalType(scale, precision);
+ }
+
+ private static int dateDays(String name, IValueMeta valueMeta, Object
valueData)
+ throws HopException {
+ long days = zoned(valueMeta, valueData).toLocalDate().toEpochDay();
+ if (days < Integer.MIN_VALUE || days > Integer.MAX_VALUE) {
+ throw new HopException(
+ "Value of field '" + name + "' does not fit in a Parquet Date (days
since the epoch)");
+ }
+ return (int) days;
+ }
+
+ private static int timeOfDayMillis(IValueMeta valueMeta, Object valueData)
throws HopException {
+ return zoned(valueMeta, valueData).get(ChronoField.MILLI_OF_DAY);
+ }
+
+ private static long timeOfDayMicros(IValueMeta valueMeta, Object valueData)
throws HopException {
+ return zoned(valueMeta, valueData).getLong(ChronoField.MICRO_OF_DAY);
+ }
+
+ /** Calendar fields of the value in the JVM time zone. */
+ private static ZonedDateTime zoned(IValueMeta valueMeta, Object valueData)
throws HopException {
+ if (valueMeta instanceof ValueMetaTimestamp timestampMeta) {
+ java.sql.Timestamp timestamp = timestampMeta.getTimestamp(valueData);
+ if (timestamp == null) {
+ throw new HopException(
+ "Unable to convert field '" + valueMeta.getName() + "' to a
timestamp");
+ }
+ return timestamp.toInstant().atZone(ZoneId.systemDefault());
+ }
+ Date date = valueMeta.getDate(valueData);
+ if (date == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to a date");
+ }
+ return Instant.ofEpochMilli(date.getTime()).atZone(ZoneId.systemDefault());
+ }
+
+ private static int int32(String name, IValueMeta valueMeta, Object valueData)
+ throws HopException {
+ long value = requireLong(valueMeta, valueData);
+ if (value < Integer.MIN_VALUE || value > Integer.MAX_VALUE) {
+ throw new HopException("Value " + value + " of field '" + name + "' does
not fit in Int32");
+ }
+ return (int) value;
+ }
+
+ private static long requireLong(IValueMeta valueMeta, Object valueData)
throws HopException {
+ Long value = valueMeta.getInteger(valueData);
+ if (value == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to an integer");
+ }
+ return value;
+ }
+
+ private static double requireDouble(IValueMeta valueMeta, Object valueData)
throws HopException {
+ Double value = valueMeta.getNumber(valueData);
+ if (value == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to a number");
+ }
+ return value;
+ }
+
+ private static boolean requireBoolean(IValueMeta valueMeta, Object valueData)
+ throws HopException {
+ Boolean value = valueMeta.getBoolean(valueData);
+ if (value == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to a boolean");
+ }
+ return value;
+ }
+
+ private static String requireString(IValueMeta valueMeta, Object valueData)
throws HopException {
+ String value = valueMeta.getString(valueData);
+ if (value == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to text");
+ }
+ return value;
+ }
+
+ private static byte[] requireBytes(IValueMeta valueMeta, Object valueData)
throws HopException {
+ byte[] value = valueMeta.getBinary(valueData);
+ if (value == null) {
+ throw new HopException("Unable to convert field '" + valueMeta.getName()
+ "' to bytes");
+ }
+ return value;
+ }
+
+ private static Integer wholeNumber(String text, String fieldName, String
what)
+ throws HopException {
+ if (Utils.isEmpty(text) || text.trim().isEmpty()) {
+ return null;
+ }
+ try {
+ return Integer.valueOf(text.trim());
+ } catch (NumberFormatException e) {
+ throw new HopException(
+ what + " '" + text + "' of field '" + fieldName + "' is not a whole
number");
+ }
+ }
+
+ private static String nameOf(ParquetField field, IValueMeta valueMeta) {
+ if (field != null && !Utils.isEmpty(field.getTargetFieldName())) {
+ return field.getTargetFieldName();
+ }
+ if (field != null && !Utils.isEmpty(field.getSourceFieldName())) {
+ return field.getSourceFieldName();
+ }
+ return valueMeta == null ? "" : valueMeta.getName();
+ }
+}
diff --git
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutput.java
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutput.java
index 951cc89bfe..541d36f5d7 100644
---
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutput.java
+++
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutput.java
@@ -246,8 +246,9 @@ public class ParquetOutput extends
BaseTransform<ParquetOutputMeta, ParquetOutpu
}
/**
- * Builds the Avro schema for the resolved output fields and converts it to
a Parquet schema.
- * Built once and reused for every split or partition file.
+ * Builds the Avro schema for the resolved output fields and converts it to
a Parquet schema. A
+ * Parquet type selected on a field replaces the column built here. Built
once and reused for
+ * every split or partition file.
*/
private MessageType buildSchema() throws HopException {
SchemaBuilder.FieldAssembler<Schema> fieldAssembler =
@@ -268,9 +269,32 @@ public class ParquetOutput extends
BaseTransform<ParquetOutputMeta, ParquetOutpu
.endUnion()
.noDefault();
}
- // Convert from Avro to Parquet schema
+ // Convert from Avro to Parquet schema, then apply a Parquet type selected
on a field.
//
- return withParquetOnlyTypes(new
AvroSchemaConverter().convert(fieldAssembler.endRecord()));
+ return applySelectedTypes(
+ withParquetOnlyTypes(new
AvroSchemaConverter().convert(fieldAssembler.endRecord())));
+ }
+
+ /**
+ * Replaces columns whose field has a Parquet type. A field without one
keeps the column built
+ * from its Hop type.
+ */
+ private MessageType applySelectedTypes(MessageType messageType) throws
HopException {
+ List<Type> types = new ArrayList<>();
+ boolean changed = false;
+ for (int i = 0; i < messageType.getFieldCount(); i++) {
+ Type type = messageType.getType(i);
+ ParquetField field = data.outputFields.get(i);
+ ParquetFieldType selected = field.parquetFieldType();
+ if (selected == null) {
+ types.add(type);
+ } else {
+ IValueMeta valueMeta =
getInputRowMeta().getValueMeta(data.sourceFieldIndexes.get(i));
+ types.add(selected.column(type.getName(), field, valueMeta));
+ changed = true;
+ }
+ }
+ return changed ? new MessageType(messageType.getName(), types) :
messageType;
}
/** The largest DECIMAL precision readers such as Spark, Hive and Trino
accept. */
@@ -392,7 +416,9 @@ public class ParquetOutput extends
BaseTransform<ParquetOutputMeta, ParquetOutpu
continue;
}
String targetFieldName = Const.NVL(field.getTargetFieldName(),
field.getSourceFieldName());
- data.outputFields.add(new ParquetField(field.getSourceFieldName(),
targetFieldName));
+ ParquetField outputField = new ParquetField(field);
+ outputField.setTargetFieldName(targetFieldName);
+ data.outputFields.add(outputField);
data.sourceFieldIndexes.add(index);
}
verifyFieldsRemain();
diff --git
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutputDialog.java
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutputDialog.java
index 7f3bc5cfc3..50669488d1 100644
---
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutputDialog.java
+++
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetOutputDialog.java
@@ -196,8 +196,29 @@ public class ParquetOutputDialog extends
BaseTransformDialog {
ColumnInfo.COLUMN_TYPE_TEXT,
false,
false),
+ new ColumnInfo(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.ParquetType.Label"),
+ ColumnInfo.COLUMN_TYPE_CCOMBO,
+ parquetTypeComboValues(),
+ true),
+ new ColumnInfo(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.Precision.Label"),
+ ColumnInfo.COLUMN_TYPE_TEXT,
+ false,
+ false),
+ new ColumnInfo(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.Scale.Label"),
+ ColumnInfo.COLUMN_TYPE_TEXT,
+ false,
+ false),
};
columns[1].setNamingSchemeType(NamingSchemeTypes.HOP_FIELD);
+ columns[2].setToolTip(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.ParquetType.Tooltip"));
+ columns[3].setToolTip(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.Precision.Tooltip"));
+ columns[4].setToolTip(
+ BaseMessages.getString(PKG,
"ParquetOutputDialog.FieldsColumn.Scale.Tooltip"));
wFields =
new TableView(
variables,
@@ -250,6 +271,9 @@ public class ParquetOutputDialog extends
BaseTransformDialog {
TableItem item = new TableItem(wFields.table, SWT.NONE);
item.setText(1, Const.NVL(field.getSourceFieldName(), ""));
item.setText(2, Const.NVL(field.getTargetFieldName(), ""));
+ item.setText(3, Const.NVL(field.getParquetType(), ""));
+ item.setText(4, Const.NVL(field.getPrecision(), ""));
+ item.setText(5, Const.NVL(field.getScale(), ""));
}
}
wFields.optimizeTableView();
@@ -271,7 +295,11 @@ public class ParquetOutputDialog extends
BaseTransformDialog {
if (wFields != null && !wFields.isDisposed()) {
List<ParquetField> fields = new ArrayList<>();
for (TableItem item : wFields.getNonEmptyItems()) {
- fields.add(new ParquetField(item.getText(1), item.getText(2)));
+ ParquetField field = new ParquetField(item.getText(1),
item.getText(2));
+ field.setParquetType(emptyToNull(item.getText(3)));
+ field.setPrecision(emptyToNull(item.getText(4)));
+ field.setScale(emptyToNull(item.getText(5)));
+ fields.add(field);
}
input.setFields(fields);
}
@@ -358,11 +386,49 @@ public class ParquetOutputDialog extends
BaseTransformDialog {
return "";
}
+ private static String[] parquetTypeComboValues() {
+ String[] codes = ParquetFieldType.codes();
+ String[] values = new String[codes.length + 1];
+ values[0] = "";
+ System.arraycopy(codes, 0, values, 1, codes.length);
+ return values;
+ }
+
+ private static String emptyToNull(String value) {
+ if (Utils.isEmpty(value)) {
+ return null;
+ }
+ String trimmed = value.trim();
+ return trimmed.isEmpty() ? null : trimmed;
+ }
+
private void getFields() {
try {
IRowMeta rowMeta = pipelineMeta.getPrevTransformFields(variables,
transformName);
BaseTransformDialog.getFieldsFromPrevious(
- rowMeta, wFields, 2, new int[] {1, 2}, new int[0], -1, -1, true,
null);
+ rowMeta,
+ wFields,
+ 2,
+ new int[] {1, 2},
+ new int[0],
+ -1,
+ -1,
+ true,
+ (tableItem, valueMeta) -> {
+ ParquetFieldType proposed =
ParquetFieldType.forValueMeta(valueMeta);
+ if (proposed != null) {
+ tableItem.setText(3, proposed.getCode());
+ }
+ if (proposed == ParquetFieldType.Decimal) {
+ if (valueMeta.getLength() > 0) {
+ tableItem.setText(4, Integer.toString(valueMeta.getLength()));
+ }
+ if (valueMeta.getPrecision() >= 0) {
+ tableItem.setText(5,
Integer.toString(valueMeta.getPrecision()));
+ }
+ }
+ return true;
+ });
} catch (Exception e) {
new ErrorDialog(
shell,
diff --git
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetWriteSupport.java
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetWriteSupport.java
index b5540718f4..cca10a0f78 100644
---
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetWriteSupport.java
+++
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/output/ParquetWriteSupport.java
@@ -51,6 +51,9 @@ public class ParquetWriteSupport extends
WriteSupport<RowMetaAndData> {
/** The logical type of the column of every field, which decides how its
values are stored. */
private final LogicalTypeAnnotation[] logicalTypes;
+ /** A Parquet type selected on the field. Null keeps the column built from
the Hop type. */
+ private final ParquetFieldType[] selectedTypes;
+
public ParquetWriteSupport(
MessageType messageType, List<Integer> sourceFieldIndexes,
List<ParquetField> fields) {
this.messageType = messageType;
@@ -60,6 +63,14 @@ public class ParquetWriteSupport extends
WriteSupport<RowMetaAndData> {
for (int i = 0; i < fields.size() && i < messageType.getFieldCount(); i++)
{
logicalTypes[i] = messageType.getType(i).getLogicalTypeAnnotation();
}
+ this.selectedTypes = new ParquetFieldType[fields.size()];
+ for (int i = 0; i < fields.size(); i++) {
+ try {
+ selectedTypes[i] = fields.get(i).parquetFieldType();
+ } catch (HopException e) {
+ throw new HopRuntimeException(e.getMessage(), e);
+ }
+ }
}
@Override
@@ -89,40 +100,48 @@ public class ParquetWriteSupport extends
WriteSupport<RowMetaAndData> {
if (!isNull) {
recordConsumer.startField(field.getTargetFieldName(), i);
- // The column type, as built by ParquetOutput.avroType(), decides
how the value is stored.
- // Anything without a column type of its own goes out as a string.
+ // A Parquet type selected on the field decides how the value is
stored. Otherwise the
+ // column type built from the Hop type does, and anything without a
column type of its
+ // own goes out as a string.
//
- LogicalTypeAnnotation logicalType = logicalTypes[i];
- switch (valueMeta.getType()) {
- case IValueMeta.TYPE_INTEGER ->
recordConsumer.addLong(valueMeta.getInteger(valueData));
- case IValueMeta.TYPE_NUMBER ->
recordConsumer.addDouble(valueMeta.getNumber(valueData));
- case IValueMeta.TYPE_BOOLEAN ->
- recordConsumer.addBoolean(valueMeta.getBoolean(valueData));
- case IValueMeta.TYPE_DATE, IValueMeta.TYPE_TIMESTAMP ->
- recordConsumer.addLong(epochValue(valueMeta, valueData,
logicalType));
- case IValueMeta.TYPE_BINARY ->
- recordConsumer.addBinary(
-
Binary.fromConstantByteArray(valueMeta.getBinary(valueData)));
- case IValueMeta.TYPE_BIGNUMBER -> {
- if (logicalType instanceof DecimalLogicalTypeAnnotation decimal)
{
- recordConsumer.addBinary(
- decimalBytes(
- field.getTargetFieldName(),
- valueMeta,
- valueMeta.getBigNumber(valueData),
- decimal));
- } else {
-
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
+ if (selectedTypes[i] != null) {
+ selectedTypes[i].write(recordConsumer, field, valueMeta,
valueData);
+ } else {
+ LogicalTypeAnnotation logicalType = logicalTypes[i];
+ switch (valueMeta.getType()) {
+ case IValueMeta.TYPE_INTEGER ->
+ recordConsumer.addLong(valueMeta.getInteger(valueData));
+ case IValueMeta.TYPE_NUMBER ->
+ recordConsumer.addDouble(valueMeta.getNumber(valueData));
+ case IValueMeta.TYPE_BOOLEAN ->
+ recordConsumer.addBoolean(valueMeta.getBoolean(valueData));
+ case IValueMeta.TYPE_DATE, IValueMeta.TYPE_TIMESTAMP ->
+ recordConsumer.addLong(epochValue(valueMeta, valueData,
logicalType));
+ case IValueMeta.TYPE_BINARY ->
+ recordConsumer.addBinary(
+
Binary.fromConstantByteArray(valueMeta.getBinary(valueData)));
+ case IValueMeta.TYPE_BIGNUMBER -> {
+ if (logicalType instanceof DecimalLogicalTypeAnnotation
decimal) {
+ recordConsumer.addBinary(
+ decimalBytes(
+ field.getTargetFieldName(),
+ valueMeta,
+ valueMeta.getBigNumber(valueData),
+ decimal));
+ } else {
+
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
+ }
}
- }
- case IValueMeta.TYPE_UUID -> {
- if (logicalType instanceof UUIDLogicalTypeAnnotation) {
-
recordConsumer.addBinary(uuidBytes(valueMeta.getString(valueData)));
- } else {
-
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
+ case IValueMeta.TYPE_UUID -> {
+ if (logicalType instanceof UUIDLogicalTypeAnnotation) {
+
recordConsumer.addBinary(uuidBytes(valueMeta.getString(valueData)));
+ } else {
+
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
+ }
}
+ default ->
+
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
}
- default ->
recordConsumer.addBinary(Binary.fromString(valueMeta.getString(valueData)));
}
recordConsumer.endField(field.getTargetFieldName(), i);
}
diff --git
a/plugins/tech/parquet/src/main/resources/org/apache/hop/parquet/transforms/output/messages/messages_en_US.properties
b/plugins/tech/parquet/src/main/resources/org/apache/hop/parquet/transforms/output/messages/messages_en_US.properties
index b4640c9d28..d7510e44c7 100644
---
a/plugins/tech/parquet/src/main/resources/org/apache/hop/parquet/transforms/output/messages/messages_en_US.properties
+++
b/plugins/tech/parquet/src/main/resources/org/apache/hop/parquet/transforms/output/messages/messages_en_US.properties
@@ -27,6 +27,12 @@ ParquetOutputDialog.DictionaryPageSize.Label=Dictionary page
size (bytes)
ParquetOutputDialog.Fields.Label=Fields (leave empty to output all input
fields)
ParquetOutputDialog.FieldsColumn.SourceField.Label=Source field
ParquetOutputDialog.FieldsColumn.TargetField.Label=Target field
+ParquetOutputDialog.FieldsColumn.ParquetType.Label=Parquet type
+ParquetOutputDialog.FieldsColumn.ParquetType.Tooltip=Logical or physical
Parquet type of this column. Leave empty to keep the type chosen from the Hop
type. Date is the calendar date in the JVM time zone, which BigQuery loads as a
DATE. TimestampMillis and TimestampMicros are instants adjusted to UTC.
Precision and scale apply to Decimal only.
+ParquetOutputDialog.FieldsColumn.Precision.Label=Precision
+ParquetOutputDialog.FieldsColumn.Precision.Tooltip=Total number of digits of a
Decimal column. Leave empty to use the length of the source field. Ignored for
every other Parquet type.
+ParquetOutputDialog.FieldsColumn.Scale.Label=Scale
+ParquetOutputDialog.FieldsColumn.Scale.Tooltip=Digits after the decimal point
of a Decimal column. Leave empty to use the precision of the source field.
Ignored for every other Parquet type.
ParquetOutputDialog.FilenameBase.Label=Base file name
ParquetOutputDialog.FilenameCompressionBeforeExtension.Label=Include
compression codec before extension
ParquetOutputDialog.FilenameCreateFolders.Label=Create parent folders
diff --git
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetFieldTest.java
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetFieldTest.java
index 672b9cbe20..7f8be5cfee 100644
---
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetFieldTest.java
+++
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetFieldTest.java
@@ -19,7 +19,10 @@ package org.apache.hop.parquet.transforms.output;
import static org.junit.jupiter.api.Assertions.assertEquals;
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 org.apache.hop.core.exception.HopException;
import org.junit.jupiter.api.Test;
/** Unit test for {@link ParquetField} */
@@ -42,9 +45,17 @@ class ParquetFieldTest {
@Test
void testCopyConstructor() {
ParquetField original = new ParquetField("source", "target");
+ original.setParquetType("Date");
+ original.setPrecision("10");
+ original.setScale("2");
ParquetField copy = new ParquetField(original);
assertEquals("source", copy.getSourceFieldName());
assertEquals("target", copy.getTargetFieldName());
+ assertEquals("Date", copy.getParquetType());
+ assertEquals("10", copy.getPrecision());
+ assertEquals("2", copy.getScale());
+ copy.setParquetType("UTF8");
+ assertEquals("Date", original.getParquetType());
}
@Test
@@ -52,7 +63,30 @@ class ParquetFieldTest {
ParquetField field = new ParquetField();
field.setSourceFieldName("in");
field.setTargetFieldName("out");
+ field.setParquetType("Int32");
+ field.setPrecision("8");
+ field.setScale("0");
assertEquals("in", field.getSourceFieldName());
assertEquals("out", field.getTargetFieldName());
+ assertEquals("Int32", field.getParquetType());
+ assertEquals("8", field.getPrecision());
+ assertEquals("0", field.getScale());
+ }
+
+ @Test
+ void testParquetType() throws Exception {
+ ParquetField field = new ParquetField("birthday", "birthday");
+ assertNull(field.parquetFieldType());
+
+ field.setParquetType(" ");
+ assertNull(field.parquetFieldType());
+
+ field.setParquetType("date");
+ assertEquals(ParquetFieldType.Date, field.parquetFieldType());
+
+ field.setParquetType("Geography");
+ HopException e = assertThrows(HopException.class, field::parquetFieldType);
+ assertTrue(e.getMessage().contains("Geography"));
+ assertTrue(e.getMessage().contains("Date"));
}
}
diff --git
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetLogicalTypeTest.java
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetLogicalTypeTest.java
new file mode 100644
index 0000000000..bce4dad9ae
--- /dev/null
+++
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetLogicalTypeTest.java
@@ -0,0 +1,475 @@
+/*
+ * 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.parquet.transforms.output;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+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 static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import java.math.BigDecimal;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.sql.Timestamp;
+import java.text.SimpleDateFormat;
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.TimeZone;
+import java.util.stream.Stream;
+import org.apache.hop.core.RowMetaAndData;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.exception.HopRuntimeException;
+import org.apache.hop.core.logging.ILoggingObject;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaBigNumber;
+import org.apache.hop.core.row.value.ValueMetaBinary;
+import org.apache.hop.core.row.value.ValueMetaBoolean;
+import org.apache.hop.core.row.value.ValueMetaDate;
+import org.apache.hop.core.row.value.ValueMetaInteger;
+import org.apache.hop.core.row.value.ValueMetaNumber;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.core.row.value.ValueMetaTimestamp;
+import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
+import org.apache.hop.pipeline.Pipeline;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transforms.mock.TransformMockHelper;
+import org.apache.parquet.hadoop.ParquetFileReader;
+import org.apache.parquet.hadoop.metadata.CompressionCodecName;
+import org.apache.parquet.io.LocalInputFile;
+import org.apache.parquet.io.api.RecordConsumer;
+import
org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
+import org.apache.parquet.schema.MessageType;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.junit.jupiter.api.io.TempDir;
+
+/** A Parquet type selected on a field is written as that logical type and
read back. */
+@ExtendWith(RestoreHopEngineEnvironmentExtension.class)
+class ParquetLogicalTypeTest {
+
+ @TempDir private Path tempDir;
+
+ @Test
+ void columnSchemas() throws Exception {
+ ParquetField field = new ParquetField("amount", "amount");
+ ValueMetaString text = new ValueMetaString("text");
+ assertEquals(
+ "optional binary text (STRING)",
+ ParquetFieldType.Utf8.column("text", field, text).toString());
+ assertEquals(
+ "optional boolean flag", ParquetFieldType.Boolean.column("flag",
field, text).toString());
+ assertEquals(
+ "optional int32 small", ParquetFieldType.Int32.column("small", field,
text).toString());
+ assertEquals(
+ "optional int64 big", ParquetFieldType.Int64.column("big", field,
text).toString());
+ assertEquals(
+ "optional float ratio", ParquetFieldType.Float.column("ratio", field,
text).toString());
+ assertEquals(
+ "optional double measure",
+ ParquetFieldType.Double.column("measure", field, text).toString());
+ assertEquals(
+ "optional binary raw", ParquetFieldType.Binary.column("raw", field,
text).toString());
+ assertEquals(
+ "optional int32 birthday (DATE)",
+ ParquetFieldType.Date.column("birthday", field, text).toString());
+ assertEquals(
+ "optional int32 clock (TIME(MILLIS,false))",
+ ParquetFieldType.TimeMillis.column("clock", field, text).toString());
+ assertEquals(
+ "optional int64 clockus (TIME(MICROS,false))",
+ ParquetFieldType.TimeMicros.column("clockus", field, text).toString());
+ assertEquals(
+ "optional int64 instant (TIMESTAMP(MILLIS,true))",
+ ParquetFieldType.TimestampMillis.column("instant", field,
text).toString());
+ assertEquals(
+ "optional int64 precise (TIMESTAMP(MICROS,true))",
+ ParquetFieldType.TimestampMicros.column("precise", field,
text).toString());
+ assertEquals(
+ "optional binary payload (JSON)",
+ ParquetFieldType.Json.column("payload", field, text).toString());
+ assertEquals(
+ "optional fixed_len_byte_array(16) id (UUID)",
+ ParquetFieldType.Uuid.column("id", field, text).toString());
+
+ field.setPrecision("10");
+ field.setScale("2");
+ assertEquals(
+ "optional binary amount (DECIMAL(10,2))",
+ ParquetFieldType.Decimal.column("amount", field, new
ValueMetaBigNumber("amount"))
+ .toString());
+ }
+
+ @Test
+ void decimalUsesExplicitPrecisionOverTheSourceField() throws Exception {
+ ParquetField field = new ParquetField("amount", "amount");
+ field.setPrecision("10");
+ field.setScale("2");
+ ValueMetaBigNumber valueMeta = new ValueMetaBigNumber("amount");
+ valueMeta.setLength(20, 5);
+
+ DecimalLogicalTypeAnnotation decimal =
ParquetFieldType.Decimal.decimal(field, valueMeta);
+ assertEquals(10, decimal.getPrecision());
+ assertEquals(2, decimal.getScale());
+
+ ParquetField fallback = new ParquetField("amount", "amount");
+ DecimalLogicalTypeAnnotation fromField =
ParquetFieldType.Decimal.decimal(fallback, valueMeta);
+ assertEquals(20, fromField.getPrecision());
+ assertEquals(5, fromField.getScale());
+ }
+
+ @Test
+ void decimalWithoutPrecisionIsRejected() {
+ ParquetField field = new ParquetField("amount", "amount");
+ HopException e =
+ assertThrows(
+ HopException.class,
+ () -> ParquetFieldType.Decimal.decimal(field, new
ValueMetaBigNumber("amount")));
+ assertTrue(e.getMessage().contains("'amount'"), e.getMessage());
+ assertTrue(e.getMessage().contains("precision"), e.getMessage());
+ }
+
+ @Test
+ void dateIsTheCalendarDateInTheJvmZone() throws Exception {
+ TimeZone original = TimeZone.getDefault();
+ TimeZone.setDefault(TimeZone.getTimeZone("Europe/Brussels"));
+ try {
+ RecordConsumer consumer = mock(RecordConsumer.class);
+ ValueMetaDate valueMeta = new ValueMetaDate("birthday");
+ // 23:30 UTC is still 1 January in Brussels, and 31 December in UTC.
+ Date date = Date.from(Instant.parse("2023-12-31T23:30:00Z"));
+ ParquetFieldType.Date.write(
+ consumer, new ParquetField("birthday", "birthday"), valueMeta, date);
+ verify(consumer).addInteger(19723);
+
+ ParquetFieldType.TimeMillis.write(
+ consumer, new ParquetField("clock", "clock"), valueMeta, date);
+ verify(consumer).addInteger(1_800_000);
+ } finally {
+ TimeZone.setDefault(original);
+ }
+ }
+
+ @Test
+ void timestampMillisIsTheUtcInstant() throws Exception {
+ RecordConsumer consumer = mock(RecordConsumer.class);
+ Date date = Date.from(Instant.parse("2024-01-01T00:00:00Z"));
+ ParquetFieldType.TimestampMillis.write(
+ consumer, new ParquetField("instant", "instant"), new
ValueMetaDate("instant"), date);
+ verify(consumer).addLong(1_704_067_200_000L);
+ }
+
+ @Test
+ void int32RejectsAValueThatDoesNotFit() {
+ HopException e =
+ assertThrows(
+ HopException.class,
+ () ->
+ ParquetFieldType.Int32.write(
+ mock(RecordConsumer.class),
+ new ParquetField("code", "code"),
+ new ValueMetaInteger("code"),
+ 5_000_000_000L));
+ assertTrue(e.getMessage().contains("'code'"), e.getMessage());
+ assertTrue(e.getMessage().contains("Int32"), e.getMessage());
+ }
+
+ @Test
+ void decimalThatDoesNotFitIsRejected() {
+ ParquetField field = new ParquetField("amount", "amount");
+ field.setPrecision("4");
+ field.setScale("2");
+ assertThrows(
+ HopRuntimeException.class,
+ () ->
+ ParquetFieldType.Decimal.write(
+ mock(RecordConsumer.class),
+ field,
+ new ValueMetaBigNumber("amount"),
+ new BigDecimal("123.45")));
+ }
+
+ @Test
+ void getFieldsProposesTimestampMillisForADateAndTheDefaultForTheOthers() {
+ assertEquals(
+ ParquetFieldType.TimestampMillis, ParquetFieldType.forValueMeta(new
ValueMetaDate("d")));
+ assertEquals(
+ ParquetFieldType.TimestampMicros,
+ ParquetFieldType.forValueMeta(new ValueMetaTimestamp("ts")));
+ assertEquals(ParquetFieldType.Utf8, ParquetFieldType.forValueMeta(new
ValueMetaString("s")));
+ assertEquals(ParquetFieldType.Int64, ParquetFieldType.forValueMeta(new
ValueMetaInteger("i")));
+ assertEquals(ParquetFieldType.Double, ParquetFieldType.forValueMeta(new
ValueMetaNumber("n")));
+ assertEquals(
+ ParquetFieldType.Boolean, ParquetFieldType.forValueMeta(new
ValueMetaBoolean("b")));
+ assertEquals(ParquetFieldType.Binary, ParquetFieldType.forValueMeta(new
ValueMetaBinary("b")));
+
+ ValueMetaBigNumber wide = new ValueMetaBigNumber("wide");
+ wide.setLength(50, 2);
+ assertEquals(ParquetFieldType.Utf8, ParquetFieldType.forValueMeta(wide));
+
+ ValueMetaBigNumber decimal = new ValueMetaBigNumber("amount");
+ decimal.setLength(10, 2);
+ assertEquals(ParquetFieldType.Decimal,
ParquetFieldType.forValueMeta(decimal));
+ assertNull(ParquetFieldType.forValueMeta(null));
+ }
+
+ @Test
+ void selectedTypesAreWrittenAndReadBack() throws Exception {
+ TimeZone original = TimeZone.getDefault();
+ TimeZone.setDefault(TimeZone.getTimeZone("Europe/Brussels"));
+ try {
+ Date birthday =
+ new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS").parse("2024-01-01
12:34:56.789");
+ Timestamp precise = Timestamp.valueOf("2024-01-01 12:34:56.123456789");
+ JsonNode payload = new ObjectMapper().readTree("{\"a\":1}");
+
+ RowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaString("text"));
+ rowMeta.addValueMeta(new ValueMetaBoolean("flag"));
+ rowMeta.addValueMeta(new ValueMetaInteger("small"));
+ rowMeta.addValueMeta(new ValueMetaInteger("big"));
+ rowMeta.addValueMeta(new ValueMetaNumber("ratio"));
+ rowMeta.addValueMeta(new ValueMetaNumber("measure"));
+ rowMeta.addValueMeta(new ValueMetaBinary("raw"));
+ rowMeta.addValueMeta(new ValueMetaDate("birthday"));
+ rowMeta.addValueMeta(new ValueMetaDate("clock"));
+ rowMeta.addValueMeta(new ValueMetaTimestamp("clockus"));
+ rowMeta.addValueMeta(new ValueMetaDate("event"));
+ rowMeta.addValueMeta(new ValueMetaDate("instant"));
+ rowMeta.addValueMeta(new ValueMetaTimestamp("precise"));
+ rowMeta.addValueMeta(new ValueMetaBigNumber("amount"));
+ rowMeta.addValueMeta(new ValueMetaString("payload"));
+ rowMeta.addValueMeta(new ValueMetaString("id"));
+
+ ParquetOutputMeta meta = new ParquetOutputMeta();
+ meta.getFields().add(typed("text", "UTF8"));
+ meta.getFields().add(typed("flag", "Boolean"));
+ meta.getFields().add(typed("small", "Int32"));
+ meta.getFields().add(typed("big", "Int64"));
+ meta.getFields().add(typed("ratio", "Float"));
+ meta.getFields().add(typed("measure", "Double"));
+ meta.getFields().add(typed("raw", "Binary"));
+ meta.getFields().add(typed("birthday", "Date"));
+ meta.getFields().add(typed("clock", "TimeMillis"));
+ meta.getFields().add(typed("clockus", "TimeMicros"));
+ // No type: a Date keeps the previous TIMESTAMP(MILLIS) mapping.
+ meta.getFields().add(new ParquetField("event", "event"));
+ meta.getFields().add(typed("instant", "TimestampMillis"));
+ meta.getFields().add(typed("precise", "TimestampMicros"));
+ ParquetField amount = typed("amount", "Decimal");
+ amount.setPrecision("10");
+ amount.setScale("2");
+ meta.getFields().add(amount);
+ meta.getFields().add(typed("payload", "JSON"));
+ meta.getFields().add(typed("id", "UUID"));
+
+ Object[] row =
+ new Object[] {
+ "h\u00e9llo",
+ true,
+ 7L,
+ 42L,
+ 1.5D,
+ 2.5D,
+ new byte[] {1, 2, 3},
+ birthday,
+ birthday,
+ precise,
+ birthday,
+ birthday,
+ precise,
+ new BigDecimal("1234.565"),
+ payload.toString(),
+ "00112233-4455-6677-8899-aabbccddeeff"
+ };
+ Path file = write(meta, rowMeta, row, new Object[rowMeta.size()]);
+
+ MessageType schema;
+ try (ParquetFileReader reader = ParquetFileReader.open(new
LocalInputFile(file))) {
+ schema = reader.getFooter().getFileMetaData().getSchema();
+ }
+ assertEquals("optional binary text (STRING)",
schema.getType("text").toString());
+ assertEquals("optional boolean flag", schema.getType("flag").toString());
+ assertEquals("optional int32 small", schema.getType("small").toString());
+ assertEquals("optional int64 big", schema.getType("big").toString());
+ assertEquals("optional float ratio", schema.getType("ratio").toString());
+ assertEquals("optional double measure",
schema.getType("measure").toString());
+ assertEquals("optional binary raw", schema.getType("raw").toString());
+ assertEquals("optional int32 birthday (DATE)",
schema.getType("birthday").toString());
+ assertEquals("optional int32 clock (TIME(MILLIS,false))",
schema.getType("clock").toString());
+ assertEquals(
+ "optional int64 clockus (TIME(MICROS,false))",
schema.getType("clockus").toString());
+ assertEquals(
+ "optional int64 event (TIMESTAMP(MILLIS,true))",
schema.getType("event").toString());
+ assertEquals(
+ "optional int64 instant (TIMESTAMP(MILLIS,true))",
schema.getType("instant").toString());
+ assertEquals(
+ "optional int64 precise (TIMESTAMP(MICROS,true))",
schema.getType("precise").toString());
+ assertEquals("optional binary amount (DECIMAL(10,2))",
schema.getType("amount").toString());
+ assertEquals("optional binary payload (JSON)",
schema.getType("payload").toString());
+ assertEquals("optional fixed_len_byte_array(16) id (UUID)",
schema.getType("id").toString());
+
+ IRowMeta readAs = ParquetTestUtil.readSchema(file.toString());
+ List<RowMetaAndData> rows =
+ ParquetTestUtil.readAllRows(file.toString(),
ParquetTestUtil.fieldsFromRowMeta(readAs));
+ RowMetaAndData read = rows.get(0);
+ assertEquals("String h\u00e9llo", render(value(read, "text")));
+ assertEquals("Boolean true", render(value(read, "flag")));
+ assertEquals("Long 7", render(value(read, "small")));
+ assertEquals("Long 42", render(value(read, "big")));
+ assertEquals("Double 1.5", render(value(read, "ratio")));
+ assertEquals("Double 2.5", render(value(read, "measure")));
+ assertEquals("bytes 010203", render(value(read, "raw")));
+ assertEquals("Date 2024-01-01 00:00:00.000", render(value(read,
"birthday")));
+ assertEquals("Timestamp 1970-01-01 12:34:56.789", render(value(read,
"clock")));
+ assertEquals("Timestamp 1970-01-01 12:34:56.123456", render(value(read,
"clockus")));
+ assertEquals(birthday.getTime(), ((Date) value(read,
"event")).getTime());
+ assertEquals(birthday.getTime(), ((Date) value(read,
"instant")).getTime());
+ assertEquals("Timestamp 2024-01-01 12:34:56.123456", render(value(read,
"precise")));
+ assertEquals("BigDecimal 1234.56", render(value(read, "amount")));
+ assertEquals("JSON {\"a\":1}", render(value(read, "payload")));
+ assertEquals("String 00112233-4455-6677-8899-aabbccddeeff",
render(value(read, "id")));
+
+ RowMetaAndData nulls = rows.get(1);
+ for (int i = 0; i < nulls.getRowMeta().size(); i++) {
+ assertNull(nulls.getData()[i], nulls.getValueMeta(i).getName());
+ }
+ } finally {
+ TimeZone.setDefault(original);
+ }
+ }
+
+ @Test
+ void unknownTypeAndMissingDecimalPrecisionFailBeforeAFileIsWritten() throws
Exception {
+ RowMeta rowMeta = new RowMeta();
+ rowMeta.addValueMeta(new ValueMetaDate("birthday"));
+ ParquetOutputMeta meta = new ParquetOutputMeta();
+ meta.getFields().add(typed("birthday", "Geography"));
+
+ HopException unknown =
+ assertThrows(HopException.class, () -> write(meta, rowMeta, new
Object[] {new Date()}));
+ assertTrue(causes(unknown).contains("Geography"), causes(unknown));
+
+ RowMeta amountRow = new RowMeta();
+ amountRow.addValueMeta(new ValueMetaBigNumber("amount"));
+ ParquetOutputMeta amountMeta = new ParquetOutputMeta();
+ amountMeta.getFields().add(typed("amount", "Decimal"));
+ HopException precision =
+ assertThrows(
+ HopException.class, () -> write(amountMeta, amountRow, new
Object[] {BigDecimal.ONE}));
+ assertTrue(causes(precision).contains("'amount'"), causes(precision));
+ assertTrue(causes(precision).contains("precision"), causes(precision));
+ }
+
+ private static ParquetField typed(String name, String parquetType) {
+ ParquetField field = new ParquetField(name, name);
+ field.setParquetType(parquetType);
+ return field;
+ }
+
+ private Path write(ParquetOutputMeta meta, IRowMeta rowMeta, Object[]...
rowsToWrite)
+ throws Exception {
+ TransformMockHelper<ParquetOutputMeta, ParquetOutputData> mockHelper =
+ new TransformMockHelper<>(
+ "Parquet Output", ParquetOutputMeta.class,
ParquetOutputData.class);
+ try {
+ when(mockHelper.logChannelFactory.create(any(),
any(ILoggingObject.class)))
+ .thenReturn(mockHelper.iLogChannel);
+ when(mockHelper.pipeline.isRunning()).thenReturn(true);
+
+ meta.setCompressionCodec(CompressionCodecName.UNCOMPRESSED);
+ meta.setFilenameIncludingCopyNr(false);
+ meta.setFilenameIncludingSplitNr(false);
+ meta.setFilenameBase(tempDir.resolve("logical").toString());
+
+ PipelineMeta pipelineMeta = new PipelineMeta();
+ TransformMeta transformMeta = new TransformMeta("Parquet Output", meta);
+ pipelineMeta.addTransform(transformMeta);
+ Pipeline pipeline = new LocalPipelineEngine(pipelineMeta);
+ ParquetOutput output =
+ spy(
+ new ParquetOutput(
+ transformMeta, meta, new ParquetOutputData(), 0,
pipelineMeta, pipeline));
+ output.setInputRowMeta(rowMeta);
+ assertTrue(output.init());
+
+ List<Object[]> remaining = new ArrayList<>(List.of(rowsToWrite));
+ doNothing().when(output).putRow(any(), any());
+ doAnswer(invocation -> remaining.isEmpty() ? null : remaining.remove(0))
+ .when(output)
+ .getRow();
+ while (output.processRow()) {
+ // keep going until the null row closes the file
+ }
+ } finally {
+ mockHelper.cleanUp();
+ }
+
+ try (Stream<Path> files = Files.list(tempDir)) {
+ return files
+ .filter(path -> path.getFileName().toString().endsWith(".parquet"))
+ .findFirst()
+ .orElseThrow();
+ }
+ }
+
+ private static Object value(RowMetaAndData row, String name) {
+ return row.getData()[row.getRowMeta().indexOfValue(name)];
+ }
+
+ private static String render(Object value) {
+ if (value instanceof Timestamp timestamp) {
+ return "Timestamp " + timestamp;
+ }
+ if (value instanceof Date date) {
+ return "Date " + new SimpleDateFormat("yyyy-MM-dd
HH:mm:ss.SSS").format(date);
+ }
+ if (value instanceof byte[] bytes) {
+ return "bytes " + java.util.HexFormat.of().formatHex(bytes);
+ }
+ if (value instanceof JsonNode json) {
+ return "JSON " + json;
+ }
+ if (value instanceof BigDecimal decimal) {
+ return "BigDecimal " + decimal.toPlainString();
+ }
+ return value.getClass().getSimpleName() + " " + value;
+ }
+
+ private static String causes(Throwable throwable) {
+ StringBuilder message = new StringBuilder();
+ while (throwable != null) {
+ message.append(throwable.getMessage()).append('\n');
+ throwable = throwable.getCause();
+ }
+ return message.toString();
+ }
+}
diff --git
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputMetaTest.java
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputMetaTest.java
index 694ed7042b..8fd7491fad 100644
---
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputMetaTest.java
+++
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputMetaTest.java
@@ -70,7 +70,11 @@ class ParquetOutputMetaTest {
original.setFilenameCompressionBeforeExtension(true);
original.setCompressionCodec(CompressionCodecName.SNAPPY);
original.setVersion(ParquetVersion.Version1);
- original.setFields(List.of(new ParquetField("id", "id_out")));
+ ParquetField id = new ParquetField("id", "id_out");
+ id.setParquetType("Date");
+ id.setPrecision("10");
+ id.setScale("2");
+ original.setFields(List.of(id));
ParquetOutputMeta copy = new ParquetOutputMeta(original);
assertEquals("/tmp/output", copy.getFilenameBase());
@@ -81,8 +85,13 @@ class ParquetOutputMetaTest {
assertEquals(ParquetVersion.Version1, copy.getVersion());
assertEquals(1, copy.getFields().size());
assertEquals("id", copy.getFields().get(0).getSourceFieldName());
+ assertEquals("Date", copy.getFields().get(0).getParquetType());
+ assertEquals("10", copy.getFields().get(0).getPrecision());
+ assertEquals("2", copy.getFields().get(0).getScale());
copy.getFields().get(0).setSourceFieldName("changed");
+ copy.getFields().get(0).setParquetType("UTF8");
assertEquals("id", original.getFields().get(0).getSourceFieldName());
+ assertEquals("Date", original.getFields().get(0).getParquetType());
}
@Test
@@ -104,7 +113,11 @@ class ParquetOutputMetaTest {
meta.setRowGroupSize("1024");
meta.setDataPageSize("512");
meta.setDictionaryPageSize("256");
- meta.getFields().add(new ParquetField("id", "id"));
+ ParquetField id = new ParquetField("id", "id");
+ id.setParquetType("Decimal");
+ id.setPrecision("10");
+ id.setScale("2");
+ meta.getFields().add(id);
meta.getFields().add(new ParquetField("name", "name"));
meta.getPartitionFields().add(new ParquetPartitionField("region"));
meta.setWriteMode(ParquetWriteMode.OverwritePartitions);
@@ -153,6 +166,11 @@ class ParquetOutputMetaTest {
assertEquals(
expected.getFields().get(i).getTargetFieldName(),
actual.getFields().get(i).getTargetFieldName());
+ assertEquals(
+ expected.getFields().get(i).getParquetType(),
actual.getFields().get(i).getParquetType());
+ assertEquals(
+ expected.getFields().get(i).getPrecision(),
actual.getFields().get(i).getPrecision());
+ assertEquals(expected.getFields().get(i).getScale(),
actual.getFields().get(i).getScale());
}
assertEquals(expected.getPartitionFields().size(),
actual.getPartitionFields().size());
for (int i = 0; i < expected.getPartitionFields().size(); i++) {
diff --git
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputTest.java
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputTest.java
index c661e1a13e..ca15418f5e 100644
---
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputTest.java
+++
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetOutputTest.java
@@ -19,6 +19,7 @@ package org.apache.hop.parquet.transforms.output;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+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 static org.mockito.ArgumentMatchers.any;
@@ -158,7 +159,11 @@ class ParquetOutputTest {
@Test
void testResolveOutputFieldsUsesConfiguredFields() throws Exception {
ParquetOutputMeta meta = new ParquetOutputMeta();
- meta.getFields().add(new ParquetField("id", "identifier"));
+ ParquetField id = new ParquetField("id", "identifier");
+ id.setParquetType("Date");
+ id.setPrecision("8");
+ id.setScale("0");
+ meta.getFields().add(id);
meta.getFields().add(new ParquetField("name", ""));
ParquetOutputData data = new ParquetOutputData();
@@ -173,7 +178,11 @@ class ParquetOutputTest {
assertEquals(2, data.outputFields.size());
assertEquals("identifier", data.outputFields.get(0).getTargetFieldName());
+ assertEquals("Date", data.outputFields.get(0).getParquetType());
+ assertEquals("8", data.outputFields.get(0).getPrecision());
+ assertEquals("0", data.outputFields.get(0).getScale());
assertEquals("name", data.outputFields.get(1).getTargetFieldName());
+ assertNull(data.outputFields.get(1).getParquetType());
}
@Test