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

commit 016f9bfdf25256ac526804ca7915d28e36d2153f
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Wed Sep 23 09:26:25 2026 +0200

    Parquet input: read drifting numeric schemas and name the column on 
conversion errors, fixes #3598
    
    - Read an int32/int64 column into a Number field and a float/double column
      into an Integer field, so files whose schemas drifted read consistently.
    - Conversion errors name the Parquet column, its type, and the Hop field.
    - ParquetStream.toString() is the plain filename, which Parquet uses in its
      own error messages.
    - The local-file fast path resolves the path through HopVfs.getFilename(),
      so a Windows drive letter is no longer lost.
    - Hop JSON fields are written as JSON-annotated columns and read back as 
JSON;
      BSON columns are read as Binary.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
---
 .../pipeline/transforms/parquet-file-input.adoc    |   5 +
 .../pipeline/transforms/parquet-file-output.adoc   |   1 +
 .../hop/parquet/transforms/input/ParquetInput.java |   2 +-
 .../parquet/transforms/input/ParquetInputMeta.java |   4 +
 .../transforms/input/ParquetRowConverter.java      |   4 +-
 .../parquet/transforms/input/ParquetStream.java    |  28 ++-
 .../transforms/input/ParquetValueConverter.java    |  56 ++++--
 .../parquet/transforms/output/ParquetOutput.java   |  40 +++-
 .../input/ParquetInputLogicalTypeTest.java         | 212 +++++++++++++++++++++
 .../transforms/input/ParquetStreamTest.java        |  88 +++++++++
 .../input/ParquetValueConverterTest.java           |  60 +++++-
 .../output/ParquetJsonRoundTripTest.java           | 162 ++++++++++++++++
 12 files changed, 636 insertions(+), 26 deletions(-)

diff --git 
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-input.adoc
 
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-input.adoc
index 879d0d549c..91df004ed9 100644
--- 
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-input.adoc
+++ 
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/parquet-file-input.adoc
@@ -38,6 +38,11 @@ Make sure to allocate enough memory to allow this.
 * Parquet Binary fields are considered to be Hop Strings but you can read them 
as Hop Binary.
 * All input values are passed to the output
 * INT96 is converted to the Hop Binary data type.
+* Columns annotated as JSON are read as the Hop JSON type, columns annotated 
as BSON as Hop Binary.
+* A column which is stored as an integer in one file and as a floating point 
number in another can
+be read as either Hop Integer or Hop Number: the value is widened or rounded 
to match the type you
+configured for the field.  This is useful when reading a set of files whose 
schemas drifted over
+time.
 
 [options="header"]
 |===
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 10b3754f9c..3d6b494ba1 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
@@ -37,6 +37,7 @@ Notes:
 * The date optionally referenced in the output file name(s) will be the start 
of the pipeline execution.
 * Hop Date and Timestamp types are serialized as milliseconds since 
`1970-01-01 00:00:00.000` UTC, with the Parquet `timestamp-millis` logical type.
 * Strings, BigNumbers, JSON and UUID values are written as UTF-8 strings.
+JSON columns also carry the Parquet JSON annotation, so that they are read 
back as JSON rather than as a String.
 * Rows are buffered in memory until a row group is full (see the *Row group 
size* option), then encoded, compressed and written.
 Memory use is therefore bounded by the row group size (times the number of 
copies and, when partitioning, the number of open partitions), not by the split 
size.
 
diff --git 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInput.java
 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInput.java
index 553fe83250..163c5d064e 100644
--- 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInput.java
+++ 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInput.java
@@ -143,7 +143,7 @@ public class ParquetInput extends 
BaseTransform<ParquetInputMeta, ParquetInputDa
         r = data.reader.read();
       }
     } catch (Exception e) {
-      throw new HopException("Error read file " + filename, e);
+      throw new HopException("Error reading Parquet file '" + filename + "'", 
e);
     } finally {
       // Every file gets its own reader; release this one before the next file 
name comes in.
       closeFile();
diff --git 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInputMeta.java
 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInputMeta.java
index 312a8fbe75..b122f437cf 100644
--- 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInputMeta.java
+++ 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetInputMeta.java
@@ -40,6 +40,7 @@ import org.apache.hop.pipeline.transform.TransformMeta;
 import org.apache.parquet.column.ColumnDescriptor;
 import org.apache.parquet.hadoop.ParquetReader;
 import org.apache.parquet.schema.LogicalTypeAnnotation;
+import 
org.apache.parquet.schema.LogicalTypeAnnotation.BsonLogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.DateLogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation;
@@ -157,6 +158,9 @@ public class ParquetInputMeta extends 
BaseTransformMeta<ParquetInput, ParquetInp
             hopType = IValueMeta.TYPE_DATE;
           } else if (logicalType instanceof JsonLogicalTypeAnnotation) {
             hopType = IValueMeta.TYPE_JSON;
+          } else if (logicalType instanceof BsonLogicalTypeAnnotation) {
+            // A BSON document is binary, reading it as text would mangle it.
+            hopType = IValueMeta.TYPE_BINARY;
           } else if (logicalType instanceof DecimalLogicalTypeAnnotation) {
             hopType = IValueMeta.TYPE_BIGNUMBER;
           } else if (logicalType instanceof IntLogicalTypeAnnotation) {
diff --git 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetRowConverter.java
 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetRowConverter.java
index 5ff78fb73c..54e8d2de3e 100644
--- 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetRowConverter.java
+++ 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetRowConverter.java
@@ -56,9 +56,7 @@ public class ParquetRowConverter extends GroupConverter {
     }
 
     return new ParquetValueConverter(
-        group,
-        rowIndex,
-        
messageType.getColumns().get(schemaIndex).getPrimitiveType().getLogicalTypeAnnotation());
+        group, rowIndex, 
messageType.getColumns().get(schemaIndex).getPrimitiveType());
   }
 
   @Override
diff --git 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetStream.java
 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetStream.java
index d5de7cbb12..a4ae097432 100644
--- 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetStream.java
+++ 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetStream.java
@@ -38,11 +38,24 @@ public class ParquetStream implements InputFile, Closeable {
   public ParquetStream(FileObject fileObject, String filename) throws 
IOException {
     this.fileObject = fileObject;
     this.filename = filename;
-    // Detect if the file is local by checking the VFS scheme
-    // For remote files, localInputFile is null. VfsSeekableInputStream will 
be used instead
-    this.isLocal = "file".equals(fileObject.getName().getScheme());
-    this.localInputFile =
-        isLocal ? new 
LocalInputFile(Paths.get(fileObject.getName().getPath())) : null;
+
+    // Use the native Parquet implementation when the file is a plain local 
file.
+    // We can't simply use the VFS path for this: on Windows it holds the 
drive letter in the
+    // root of the file name, so the path alone would resolve against the 
current drive.
+    // HopVfs.getFilename() puts the two back together and hands back a URI 
for anything it
+    // can't express as a local path, a Windows network share for example.
+    // For those, and for remote files, localInputFile is null and 
VfsSeekableInputStream is
+    // used instead.
+    //
+    String localFilename = null;
+    if ("file".equals(fileObject.getName().getScheme())) {
+      String candidate = HopVfs.getFilename(fileObject);
+      if (!candidate.startsWith("file:")) {
+        localFilename = candidate;
+      }
+    }
+    this.isLocal = localFilename != null;
+    this.localInputFile = isLocal ? new 
LocalInputFile(Paths.get(localFilename)) : null;
   }
 
   @Override
@@ -83,7 +96,10 @@ public class ParquetStream implements InputFile, Closeable {
 
   @Override
   public String toString() {
-    return "ParquetStream of file '" + filename + "'";
+    // Parquet uses this to identify the file in its error messages ("...in 
file %s"), so keep
+    // it to the plain file name.
+    //
+    return filename;
   }
 
   // SeekableInputStream implementation for remote files
diff --git 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetValueConverter.java
 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetValueConverter.java
index ddae88973c..c3e2b31761 100644
--- 
a/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetValueConverter.java
+++ 
b/plugins/tech/parquet/src/main/java/org/apache/hop/parquet/transforms/input/ParquetValueConverter.java
@@ -36,20 +36,46 @@ import org.apache.parquet.io.api.PrimitiveConverter;
 import org.apache.parquet.schema.LogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.DateLogicalTypeAnnotation;
 import 
org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation;
+import org.apache.parquet.schema.PrimitiveType;
 
 public class ParquetValueConverter extends PrimitiveConverter {
 
   private final RowMetaAndData group;
   private final IValueMeta valueMeta;
   private final int rowIndex;
+  private final PrimitiveType primitiveType;
   private final LogicalTypeAnnotation logicalTypeAnnotation;
 
-  public ParquetValueConverter(
-      RowMetaAndData group, int rowIndex, LogicalTypeAnnotation 
logicalTypeAnnotation) {
+  public ParquetValueConverter(RowMetaAndData group, int rowIndex, 
PrimitiveType primitiveType) {
     this.group = group;
     this.valueMeta = group.getValueMeta(rowIndex);
     this.rowIndex = rowIndex;
-    this.logicalTypeAnnotation = logicalTypeAnnotation;
+    this.primitiveType = primitiveType;
+    // A column which isn't mapped to a field is never converted, so it is 
allowed to come in
+    // without a type.
+    this.logicalTypeAnnotation =
+        primitiveType == null ? null : 
primitiveType.getLogicalTypeAnnotation();
+  }
+
+  /**
+   * Build an error which names both ends of the conversion: the Parquet 
column with its physical
+   * type and the Hop field with its type. Files read by a single transform 
all have to match the
+   * configured fields, so a mismatch here usually means the files don't share 
the same schema.
+   *
+   * @return the exception to throw
+   */
+  private HopRuntimeException conversionError() {
+    return new HopRuntimeException(
+        "Unable to convert Parquet column '"
+            + primitiveType.getName()
+            + "' of type "
+            + primitiveType.getPrimitiveTypeName()
+            + (logicalTypeAnnotation == null ? "" : " (" + 
logicalTypeAnnotation + ")")
+            + " to field '"
+            + valueMeta.getName()
+            + "' of type "
+            + valueMeta.getTypeDesc()
+            + ". Please verify that all the files being read have the same 
schema.");
   }
 
   @Override
@@ -122,8 +148,7 @@ public class ParquetValueConverter extends 
PrimitiveConverter {
           break;
         }
       default:
-        throw new HopRuntimeException(
-            "Unable to convert Binary source data to type " + 
valueMeta.getTypeDesc());
+        throw conversionError();
     }
     group.getData()[rowIndex] = object;
   }
@@ -138,6 +163,11 @@ public class ParquetValueConverter extends 
PrimitiveConverter {
       case IValueMeta.TYPE_INTEGER:
         object = value;
         break;
+      case IValueMeta.TYPE_NUMBER:
+        // An integer column read into a Number field: the same column can be 
stored as an
+        // integer in one file and as a double in the next one.
+        object = (double) value;
+        break;
       case IValueMeta.TYPE_STRING:
         object = Long.toString(value);
         break;
@@ -162,8 +192,7 @@ public class ParquetValueConverter extends 
PrimitiveConverter {
         object = convertToTimestamp(value, this.logicalTypeAnnotation);
         break;
       default:
-        throw new HopRuntimeException(
-            "Unable to convert Long source data to type " + 
valueMeta.getTypeDesc());
+        throw conversionError();
     }
     group.getData()[rowIndex] = object;
   }
@@ -173,14 +202,17 @@ public class ParquetValueConverter extends 
PrimitiveConverter {
     if (rowIndex < 0) {
       return;
     }
+    // A floating point column can be read into an Integer field: the same 
column can be stored
+    // as a double in one file and as an integer in the next one. We round it 
like Hop does
+    // everywhere else when converting a Number to an Integer.
+    //
     Object object =
         switch (valueMeta.getType()) {
           case IValueMeta.TYPE_NUMBER -> value;
+          case IValueMeta.TYPE_INTEGER -> Math.round(value);
           case IValueMeta.TYPE_STRING -> Double.toString(value);
           case IValueMeta.TYPE_BIGNUMBER -> BigDecimal.valueOf(value);
-          default ->
-              throw new HopRuntimeException(
-                  "Unable to convert Double/Float source data to type " + 
valueMeta.getTypeDesc());
+          default -> throw conversionError();
         };
     group.getData()[rowIndex] = object;
   }
@@ -195,9 +227,7 @@ public class ParquetValueConverter extends 
PrimitiveConverter {
           case IValueMeta.TYPE_BOOLEAN -> value;
           case IValueMeta.TYPE_STRING -> value ? "true" : "false";
           case IValueMeta.TYPE_INTEGER -> value ? 1L : 0L;
-          default ->
-              throw new HopRuntimeException(
-                  "Unable to convert Boolean source data to type " + 
valueMeta.getTypeDesc());
+          default -> throw conversionError();
         };
     group.getData()[rowIndex] = object;
   }
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 4f83f98ecd..6811b63cd1 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
@@ -28,6 +28,7 @@ import java.util.Iterator;
 import java.util.List;
 import java.util.Locale;
 import java.util.Map;
+import java.util.Set;
 import java.util.UUID;
 import org.apache.avro.LogicalTypes;
 import org.apache.avro.Schema;
@@ -52,7 +53,10 @@ import org.apache.parquet.avro.AvroSchemaConverter;
 import org.apache.parquet.column.ParquetProperties;
 import org.apache.parquet.hadoop.ParquetFileWriter;
 import org.apache.parquet.hadoop.ParquetWriter;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
 import org.apache.parquet.schema.MessageType;
+import org.apache.parquet.schema.Type;
+import org.apache.parquet.schema.Types;
 
 public class ParquetOutput extends BaseTransform<ParquetOutputMeta, 
ParquetOutputData> {
 
@@ -265,7 +269,7 @@ public class ParquetOutput extends 
BaseTransform<ParquetOutputMeta, ParquetOutpu
     }
     // Convert from Avro to Parquet schema
     //
-    return new AvroSchemaConverter().convert(fieldAssembler.endRecord());
+    return annotateJsonFields(new 
AvroSchemaConverter().convert(fieldAssembler.endRecord()));
   }
 
   /**
@@ -294,6 +298,40 @@ public class ParquetOutput extends 
BaseTransform<ParquetOutputMeta, ParquetOutpu
     };
   }
 
+  /**
+   * Avro has no JSON logical type, so the schema coming out of the Avro 
converter describes a Hop
+   * JSON field as a plain string. Annotate those columns as JSON again so 
that readers, Parquet
+   * Input included, recognise them as JSON instead of String.
+   *
+   * @param messageType the schema as converted from Avro
+   * @return the same schema with the JSON columns annotated
+   */
+  private MessageType annotateJsonFields(MessageType messageType) {
+    Set<String> jsonFieldNames = new HashSet<>();
+    for (int i = 0; i < data.outputFields.size(); i++) {
+      IValueMeta valueMeta = 
getInputRowMeta().getValueMeta(data.sourceFieldIndexes.get(i));
+      if (valueMeta.getType() == IValueMeta.TYPE_JSON) {
+        jsonFieldNames.add(data.outputFields.get(i).getTargetFieldName());
+      }
+    }
+    if (jsonFieldNames.isEmpty()) {
+      return messageType;
+    }
+
+    List<Type> types = new ArrayList<>();
+    for (Type type : messageType.getFields()) {
+      if (jsonFieldNames.contains(type.getName()) && type.isPrimitive()) {
+        types.add(
+            Types.primitive(type.asPrimitiveType().getPrimitiveTypeName(), 
type.getRepetition())
+                .as(LogicalTypeAnnotation.jsonType())
+                .named(type.getName()));
+      } else {
+        types.add(type);
+      }
+    }
+    return new MessageType(messageType.getName(), types);
+  }
+
   void resolveOutputFields() throws HopException {
     data.outputFields = new ArrayList<>();
     data.sourceFieldIndexes = new ArrayList<>();
diff --git 
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetInputLogicalTypeTest.java
 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetInputLogicalTypeTest.java
new file mode 100644
index 0000000000..ff576e45e0
--- /dev/null
+++ 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetInputLogicalTypeTest.java
@@ -0,0 +1,212 @@
+/*
+ * 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.input;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+import java.math.BigDecimal;
+import java.nio.file.Path;
+import java.sql.Timestamp;
+import java.time.LocalDate;
+import java.time.ZoneId;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import org.apache.commons.vfs2.FileObject;
+import org.apache.hop.core.HopClientEnvironment;
+import org.apache.hop.core.RowMetaAndData;
+import org.apache.hop.core.vfs.HopVfs;
+import org.apache.parquet.example.data.Group;
+import org.apache.parquet.example.data.simple.SimpleGroupFactory;
+import org.apache.parquet.hadoop.ParquetReader;
+import org.apache.parquet.hadoop.ParquetWriter;
+import org.apache.parquet.hadoop.example.ExampleParquetWriter;
+import org.apache.parquet.io.LocalOutputFile;
+import org.apache.parquet.io.api.Binary;
+import org.apache.parquet.schema.LogicalTypeAnnotation;
+import org.apache.parquet.schema.MessageType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Types;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/**
+ * Reads a file holding annotated Parquet columns back through the real read 
path, to guard the
+ * logical type handling in {@link ParquetValueConverter}.
+ */
+class ParquetInputLogicalTypeTest {
+
+  /** 2024-01-31, the day the schema mismatch behind issue #3598 was reported. 
*/
+  private static final LocalDate DATE = LocalDate.of(2024, 1, 31);
+
+  private static final long EPOCH_MILLIS = 1706716800123L;
+
+  private static final MessageType SCHEMA =
+      Types.buildMessage()
+          .optional(PrimitiveTypeName.INT32)
+          .as(LogicalTypeAnnotation.dateType())
+          .named("date_field")
+          .optional(PrimitiveTypeName.INT64)
+          .as(LogicalTypeAnnotation.timestampType(true, 
LogicalTypeAnnotation.TimeUnit.MILLIS))
+          .named("timestamp_millis_field")
+          .optional(PrimitiveTypeName.INT64)
+          .as(LogicalTypeAnnotation.timestampType(true, 
LogicalTypeAnnotation.TimeUnit.MICROS))
+          .named("timestamp_micros_field")
+          .optional(PrimitiveTypeName.INT64)
+          .as(LogicalTypeAnnotation.decimalType(2, 18))
+          .named("decimal_field")
+          .optional(PrimitiveTypeName.BINARY)
+          .as(LogicalTypeAnnotation.jsonType())
+          .named("json_field")
+          .optional(PrimitiveTypeName.BINARY)
+          .as(LogicalTypeAnnotation.stringType())
+          .named("string_field")
+          .optional(PrimitiveTypeName.INT64)
+          .named("long_field")
+          .optional(PrimitiveTypeName.DOUBLE)
+          .named("double_field")
+          .named("LogicalTypes");
+
+  @BeforeAll
+  static void setUpBeforeAll() throws Exception {
+    HopClientEnvironment.init();
+  }
+
+  private static ParquetField field(String name, String type) {
+    return new ParquetField(name, name, type, null, "-1", "-1");
+  }
+
+  private static Path writeFile(Path folder) throws Exception {
+    Path file = folder.resolve("logical-types.parquet");
+    SimpleGroupFactory groupFactory = new SimpleGroupFactory(SCHEMA);
+
+    try (ParquetWriter<Group> writer =
+        ExampleParquetWriter.builder(new 
LocalOutputFile(file)).withType(SCHEMA).build()) {
+      writer.write(
+          groupFactory
+              .newGroup()
+              .append("date_field", (int) DATE.toEpochDay())
+              .append("timestamp_millis_field", EPOCH_MILLIS)
+              .append("timestamp_micros_field", EPOCH_MILLIS * 1000L)
+              .append("decimal_field", 123456L)
+              .append("json_field", Binary.fromString("{\"a\":1}"))
+              .append("string_field", Binary.fromString("hop"))
+              .append("long_field", 9081496L)
+              .append("double_field", 9081496.6d));
+    }
+    return file;
+  }
+
+  private static RowMetaAndData read(Path file, List<ParquetField> fields) 
throws Exception {
+    FileObject fileObject = HopVfs.getFileObject(file.toString());
+    try (ParquetStream stream = new ParquetStream(fileObject, file.toString());
+        ParquetReader<RowMetaAndData> reader =
+            new ParquetReaderBuilder<>(new ParquetReadSupport(fields), 
stream).build()) {
+      return reader.read();
+    }
+  }
+
+  /** Every annotated column still has to be converted using its logical type. 
*/
+  @Test
+  void testLogicalTypesAreHonoured(@TempDir Path folder) throws Exception {
+    Path file = writeFile(folder);
+
+    List<ParquetField> fields = new ArrayList<>();
+    fields.add(field("date_field", "Date"));
+    fields.add(field("timestamp_millis_field", "Timestamp"));
+    fields.add(field("timestamp_micros_field", "Timestamp"));
+    fields.add(field("decimal_field", "BigNumber"));
+    fields.add(field("json_field", "JSON"));
+    fields.add(field("string_field", "String"));
+
+    RowMetaAndData row = read(file, fields);
+    assertNotNull(row);
+
+    // The DATE annotation makes this an epoch day rather than a number of 
millis.
+    Date expectedDate = 
Date.from(DATE.atStartOfDay(ZoneId.systemDefault()).toInstant());
+    assertEquals(expectedDate, row.getData()[0]);
+
+    // Both timestamps describe the same instant, in a different unit.
+    assertEquals(new Timestamp(EPOCH_MILLIS), row.getData()[1]);
+    assertEquals(new Timestamp(EPOCH_MILLIS), row.getData()[2]);
+
+    // The DECIMAL annotation carries the scale: 123456 with scale 2 is 1234.56
+    assertEquals(0, new BigDecimal("1234.56").compareTo((BigDecimal) 
row.getData()[3]));
+
+    assertEquals("{\"a\":1}", row.getData()[4].toString());
+    assertEquals("hop", row.getData()[5]);
+  }
+
+  /**
+   * A DATE column read without its logical type would come out as an instant 
near the epoch. This
+   * pins the annotation actually reaching the converter.
+   */
+  @Test
+  void testDateColumnIsNotReadAsMillis(@TempDir Path folder) throws Exception {
+    Path file = writeFile(folder);
+
+    RowMetaAndData row = read(file, List.of(field("date_field", "Date")));
+
+    assertEquals(
+        DATE, ((Date) 
row.getData()[0]).toInstant().atZone(ZoneId.systemDefault()).toLocalDate());
+  }
+
+  /** The widening conversions added for issue #3598, through the real read 
path. */
+  @Test
+  void testNumericTypeMismatchIsWidened(@TempDir Path folder) throws Exception 
{
+    Path file = writeFile(folder);
+
+    // An int64 column read into a Number field and a double column read into 
an Integer field.
+    RowMetaAndData row =
+        read(file, List.of(field("long_field", "Number"), 
field("double_field", "Integer")));
+
+    assertEquals(9081496.0, (Double) row.getData()[0], 0.0);
+    assertEquals(9081497L, row.getData()[1]);
+  }
+
+  /** Columns which aren't requested stay out of the row. */
+  @Test
+  void testUnrequestedColumnsAreNotRead(@TempDir Path folder) throws Exception 
{
+    Path file = writeFile(folder);
+
+    RowMetaAndData row = read(file, List.of(field("string_field", "String")));
+
+    assertEquals(1, row.getRowMeta().size());
+    assertEquals("hop", row.getData()[0]);
+  }
+
+  /** Nulls stay null rather than being converted. */
+  @Test
+  void testNullValues(@TempDir Path folder) throws Exception {
+    Path file = folder.resolve("nulls.parquet");
+    SimpleGroupFactory groupFactory = new SimpleGroupFactory(SCHEMA);
+    try (ParquetWriter<Group> writer =
+        ExampleParquetWriter.builder(new 
LocalOutputFile(file)).withType(SCHEMA).build()) {
+      writer.write(groupFactory.newGroup().append("string_field", 
Binary.fromString("hop")));
+    }
+
+    RowMetaAndData row =
+        read(file, List.of(field("date_field", "Date"), field("string_field", 
"String")));
+
+    assertNull(row.getData()[0]);
+    assertEquals("hop", row.getData()[1]);
+  }
+}
diff --git 
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetStreamTest.java
 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetStreamTest.java
new file mode 100644
index 0000000000..7e474cd9aa
--- /dev/null
+++ 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetStreamTest.java
@@ -0,0 +1,88 @@
+/*
+ * 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.input;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.io.OutputStream;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Path;
+import org.apache.commons.vfs2.FileObject;
+import org.apache.hop.core.HopClientEnvironment;
+import org.apache.hop.core.vfs.HopVfs;
+import org.apache.parquet.io.SeekableInputStream;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/** Unit test for {@link ParquetStream} */
+class ParquetStreamTest {
+
+  private static final byte[] CONTENT = "Apache Hop 
Parquet".getBytes(StandardCharsets.UTF_8);
+
+  @BeforeAll
+  static void setUpBeforeAll() throws Exception {
+    HopClientEnvironment.init();
+  }
+
+  private static FileObject createFile(Path folder) throws Exception {
+    FileObject fileObject = 
HopVfs.getFileObject(folder.resolve("test.parquet").toString());
+    try (OutputStream outputStream = HopVfs.getOutputStream(fileObject, 
false)) {
+      outputStream.write(CONTENT);
+    }
+    return fileObject;
+  }
+
+  /**
+   * A local file is read through the native Parquet implementation. That one 
needs a real local
+   * path: on Windows the drive letter lives in the root of the VFS file name, 
so using the VFS path
+   * on its own would resolve against the current drive.
+   */
+  @Test
+  void testLocalFileLength(@TempDir Path folder) throws Exception {
+    FileObject fileObject = createFile(folder);
+
+    try (ParquetStream stream = new ParquetStream(fileObject, 
fileObject.toString())) {
+      assertEquals(CONTENT.length, stream.getLength());
+    }
+  }
+
+  @Test
+  void testLocalFileRead(@TempDir Path folder) throws Exception {
+    FileObject fileObject = createFile(folder);
+
+    try (ParquetStream stream = new ParquetStream(fileObject, 
fileObject.toString());
+        SeekableInputStream inputStream = stream.newStream()) {
+      byte[] buffer = new byte[CONTENT.length];
+      inputStream.readFully(buffer);
+      assertEquals(
+          new String(CONTENT, StandardCharsets.UTF_8), new String(buffer, 
StandardCharsets.UTF_8));
+    }
+  }
+
+  /** Parquet puts this in its error messages ("...in file %s") so it has to 
stay readable. */
+  @Test
+  void testToStringIsThePlainFilename(@TempDir Path folder) throws Exception {
+    FileObject fileObject = createFile(folder);
+    String filename = folder.resolve("test.parquet").toString();
+
+    try (ParquetStream stream = new ParquetStream(fileObject, filename)) {
+      assertEquals(filename, stream.toString());
+    }
+  }
+}
diff --git 
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetValueConverterTest.java
 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetValueConverterTest.java
index 73dfe54c38..fc0aa940e9 100644
--- 
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetValueConverterTest.java
+++ 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/input/ParquetValueConverterTest.java
@@ -21,6 +21,7 @@ import static 
org.junit.jupiter.api.Assertions.assertArrayEquals;
 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 com.fasterxml.jackson.databind.JsonNode;
 import java.math.BigDecimal;
@@ -31,6 +32,7 @@ import java.time.LocalDate;
 import java.time.ZoneId;
 import java.util.Date;
 import java.util.TimeZone;
+import org.apache.hop.core.HopClientEnvironment;
 import org.apache.hop.core.RowMetaAndData;
 import org.apache.hop.core.exception.HopRuntimeException;
 import org.apache.hop.core.row.IValueMeta;
@@ -47,6 +49,10 @@ import org.apache.hop.core.row.value.ValueMetaTimestamp;
 import org.apache.parquet.io.api.Binary;
 import org.apache.parquet.schema.LogicalTypeAnnotation;
 import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit;
+import org.apache.parquet.schema.PrimitiveType;
+import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
+import org.apache.parquet.schema.Type;
+import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.Test;
 
 /** Unit test for {@link ParquetValueConverter}: every Parquet primitive into 
every Hop type. */
@@ -54,11 +60,25 @@ class ParquetValueConverterTest {
 
   private RowMetaAndData group;
 
+  @BeforeAll
+  static void setUpBeforeAll() throws Exception {
+    // The value types need to be registered for IValueMeta.getTypeDesc() to 
name them.
+    HopClientEnvironment.init();
+  }
+
   private ParquetValueConverter converter(IValueMeta valueMeta, 
LogicalTypeAnnotation annotation) {
+    return converter(valueMeta, PrimitiveTypeName.BINARY, annotation);
+  }
+
+  private ParquetValueConverter converter(
+      IValueMeta valueMeta, PrimitiveTypeName typeName, LogicalTypeAnnotation 
annotation) {
     RowMeta rowMeta = new RowMeta();
     rowMeta.addValueMeta(valueMeta);
     group = new RowMetaAndData(rowMeta, new Object[1]);
-    return new ParquetValueConverter(group, 0, annotation);
+    PrimitiveType column =
+        new PrimitiveType(Type.Repetition.OPTIONAL, typeName, "column")
+            .withLogicalTypeAnnotation(annotation);
+    return new ParquetValueConverter(group, 0, column);
   }
 
   private Object value() {
@@ -241,7 +261,14 @@ class ParquetValueConverterTest {
     converter(new ValueMetaNumber("n"), null).addFloat(2.5F);
     assertEquals(2.5D, value());
 
-    ParquetValueConverter c = converter(new ValueMetaInteger("i"), null);
+    // The same column can be a double in one file and an integer in the next 
one (#3598).
+    converter(new ValueMetaInteger("i"), null).addDouble(1.5D);
+    assertEquals(2L, value());
+
+    converter(new ValueMetaInteger("i"), null).addFloat(41.4F);
+    assertEquals(41L, value());
+
+    ParquetValueConverter c = converter(new ValueMetaBoolean("b"), null);
     assertThrows(HopRuntimeException.class, () -> c.addDouble(1.5D));
   }
 
@@ -289,4 +316,33 @@ class ParquetValueConverterTest {
             Binary.fromConstantByteArray(new byte[] {(byte) 0xFF}), 38, 1);
     assertEquals(0, new BigDecimal("-0.1").compareTo(large));
   }
+
+  /** The same column can be an integer in one file and a double in the next 
one (#3598). */
+  @Test
+  void longAndIntIntoNumber() {
+    converter(new ValueMetaNumber("n"), PrimitiveTypeName.INT64, 
null).addLong(9081496L);
+    assertEquals(9081496.0D, value());
+
+    converter(new ValueMetaNumber("n"), PrimitiveTypeName.INT32, 
null).addInt(42);
+    assertEquals(42.0D, value());
+  }
+
+  /**
+   * A conversion which really can't be made names the Parquet column, its 
type, the Hop field and
+   * its type, so the schema mismatch behind it can be found (#3598).
+   */
+  @Test
+  void unsupportedConversionNamesColumnAndField() {
+    ParquetValueConverter c =
+        converter(new ValueMetaBoolean("flag"), PrimitiveTypeName.DOUBLE, 
null);
+
+    HopRuntimeException e = assertThrows(HopRuntimeException.class, () -> 
c.addDouble(1.5D));
+
+    String message = e.getMessage();
+    assertTrue(message.contains("'column'"), message);
+    assertTrue(message.contains("DOUBLE"), message);
+    assertTrue(message.contains("'flag'"), message);
+    assertTrue(message.contains("Boolean"), message);
+    assertTrue(message.contains("same schema"), message);
+  }
 }
diff --git 
a/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetJsonRoundTripTest.java
 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetJsonRoundTripTest.java
new file mode 100644
index 0000000000..85c777aedc
--- /dev/null
+++ 
b/plugins/tech/parquet/src/test/java/org/apache/hop/parquet/transforms/output/ParquetJsonRoundTripTest.java
@@ -0,0 +1,162 @@
+/*
+ * 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.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.when;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.List;
+import java.util.stream.Stream;
+import org.apache.hop.core.RowMetaAndData;
+import org.apache.hop.core.logging.ILoggingObject;
+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.ValueMetaJson;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.junit.rules.RestoreHopEngineEnvironmentExtension;
+import org.apache.hop.parquet.transforms.input.ParquetField;
+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.metadata.CompressionCodecName;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.junit.jupiter.api.io.TempDir;
+
+/**
+ * Avro, which the output schema is built with, has no JSON logical type. 
These tests pin that a Hop
+ * JSON field still comes back out of the file as a JSON field rather than as 
a String.
+ */
+@ExtendWith(RestoreHopEngineEnvironmentExtension.class)
+class ParquetJsonRoundTripTest {
+
+  private static final String JSON = "{\"id\":1,\"tags\":[\"a\",\"b\"]}";
+
+  @TempDir private Path tempDir;
+
+  private TransformMockHelper<ParquetOutputMeta, ParquetOutputData> mockHelper;
+
+  @BeforeEach
+  void setUp() {
+    mockHelper =
+        new TransformMockHelper<>(
+            "Parquet Output", ParquetOutputMeta.class, 
ParquetOutputData.class);
+    when(mockHelper.logChannelFactory.create(any(), any(ILoggingObject.class)))
+        .thenReturn(mockHelper.iLogChannel);
+    when(mockHelper.pipeline.isRunning()).thenReturn(true);
+  }
+
+  @AfterEach
+  void tearDown() {
+    mockHelper.cleanUp();
+  }
+
+  /** A Hop JSON field has to be annotated as JSON in the schema of the file 
we write. */
+  @Test
+  void testJsonFieldIsReadBackAsJson() throws Exception {
+    Path file = writeOneRow();
+
+    IRowMeta rowMeta = ParquetTestUtil.readSchema(file.toString());
+
+    assertEquals(IValueMeta.TYPE_JSON, 
rowMeta.getValueMeta(rowMeta.indexOfValue("doc")).getType());
+    // A plain string field has to stay a String, so the annotation isn't 
applied to everything.
+    assertEquals(
+        IValueMeta.TYPE_STRING, 
rowMeta.getValueMeta(rowMeta.indexOfValue("name")).getType());
+  }
+
+  /** And the value itself has to survive the round trip. */
+  @Test
+  void testJsonValueSurvivesTheRoundTrip() throws Exception {
+    Path file = writeOneRow();
+
+    List<ParquetField> fields =
+        List.of(
+            new ParquetField("doc", "doc", "JSON", null, "-1", "-1"),
+            new ParquetField("name", "name", "String", null, "-1", "-1"));
+    List<RowMetaAndData> rows = ParquetTestUtil.readAllRows(file.toString(), 
fields);
+
+    assertEquals(1, rows.size());
+    RowMetaAndData row = rows.get(0);
+    assertEquals(IValueMeta.TYPE_JSON, row.getValueMeta(0).getType());
+    assertEquals(
+        new ObjectMapper().readTree(JSON), new 
ObjectMapper().readTree(row.getString(0, "")));
+    assertEquals("hop", row.getString(1, ""));
+  }
+
+  private Path writeOneRow() throws Exception {
+    ParquetOutputMeta meta = new ParquetOutputMeta();
+    meta.setCompressionCodec(CompressionCodecName.UNCOMPRESSED);
+    meta.setFilenameIncludingSplitNr(false);
+    meta.setFilenameBase(tempDir.resolve("docs").toString());
+    meta.setRowGroupSize("4096");
+    meta.setDataPageSize("1024");
+    meta.setDictionaryPageSize("512");
+
+    ParquetOutputData data = new ParquetOutputData();
+    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, data, 0, pipelineMeta, 
pipeline));
+
+    RowMeta rowMeta = new RowMeta();
+    rowMeta.addValueMeta(new ValueMetaJson("doc"));
+    rowMeta.addValueMeta(new ValueMetaString("name"));
+    output.setInputRowMeta(rowMeta);
+    assertTrue(output.init());
+
+    // A single row, spelled out rather than List.of() so it stays a list of 
one Object[].
+    List<Object[]> remaining = new ArrayList<>();
+    remaining.add(new Object[] {new ObjectMapper().readTree(JSON), "hop"});
+
+    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
+    }
+
+    return onlyFile(tempDir);
+  }
+
+  private static Path onlyFile(Path folder) throws IOException {
+    try (Stream<Path> stream = Files.list(folder)) {
+      List<Path> files =
+          
stream.filter(Files::isRegularFile).sorted(Comparator.naturalOrder()).toList();
+      assertEquals(1, files.size(), () -> "Expected one file in " + folder + " 
but got " + files);
+      return files.get(0);
+    }
+  }
+}

Reply via email to