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