diff --git a/java/core/src/java/org/apache/orc/impl/RecordReaderImpl.java b/java/core/src/java/org/apache/orc/impl/RecordReaderImpl.java index 34ffe45204..4ca0702872 100644 --- a/java/core/src/java/org/apache/orc/impl/RecordReaderImpl.java +++ b/java/core/src/java/org/apache/orc/impl/RecordReaderImpl.java @@ -1311,7 +1311,9 @@ private void clearStreams() { */ private void readStripe() throws IOException { StripeInformation stripe = beginReadStripe(); - planner.parseStripe(stripe, fileIncluded); + // Reuse the footer that beginReadStripe already read to avoid a + // duplicate readStripeFooter call inside StripePlanner.parseStripe. + planner.parseStripe(stripe, fileIncluded, stripeFooter); includedRowGroups = pickRowGroups(); // move forward to the first unskipped row diff --git a/java/core/src/java/org/apache/orc/impl/reader/StripePlanner.java b/java/core/src/java/org/apache/orc/impl/reader/StripePlanner.java index 0ec93611cd..18c9a270b6 100644 --- a/java/core/src/java/org/apache/orc/impl/reader/StripePlanner.java +++ b/java/core/src/java/org/apache/orc/impl/reader/StripePlanner.java @@ -140,7 +140,22 @@ public StripePlanner(StripePlanner old) { */ public StripePlanner parseStripe(StripeInformation stripe, boolean[] columnInclude) throws IOException { - OrcProto.StripeFooter footer = dataReader.readStripeFooter(stripe); + return parseStripe(stripe, columnInclude, dataReader.readStripeFooter(stripe)); + } + + /** + * Parse a new stripe using a footer that was already read. Resets the + * current stripe state. This avoids a redundant {@code readStripeFooter} + * call when the caller (e.g. {@code RecordReaderImpl}) has already read + * the footer for the same stripe. + * @param stripe the new stripe + * @param columnInclude an array with true for each column to read + * @param footer the stripe footer that was already read for {@code stripe} + * @return this for method chaining + */ + public StripePlanner parseStripe(StripeInformation stripe, + boolean[] columnInclude, + OrcProto.StripeFooter footer) throws IOException { currentStripeId = stripe.getStripeId(); originalStripeId = stripe.getEncryptionStripeId(); writerTimezone = footer.getWriterTimezone(); diff --git a/java/core/src/test/org/apache/orc/TestMinSeekSize.java b/java/core/src/test/org/apache/orc/TestMinSeekSize.java index 0040501fc4..169bce1ae8 100644 --- a/java/core/src/test/org/apache/orc/TestMinSeekSize.java +++ b/java/core/src/test/org/apache/orc/TestMinSeekSize.java @@ -183,8 +183,9 @@ public void readAlternateColumnsWMinSeekSize() throws IOException { FileSystem.Statistics stats = readEnd(); double p = readPercentage(stats, fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - // Read all bytes - assertTrue(p >= 5.9); + // Read all bytes. Threshold lowered from 5.9: footer dedup removes one + // footer read per stripe, reducing total bytes read. + assertTrue(p >= 5.0); } private double readPercentage(FileSystem.Statistics stats, long fileSize) { diff --git a/java/core/src/test/org/apache/orc/TestRowFilteringComplexTypesNulls.java b/java/core/src/test/org/apache/orc/TestRowFilteringComplexTypesNulls.java index 248e6c88dd..fde86eb19e 100644 --- a/java/core/src/test/org/apache/orc/TestRowFilteringComplexTypesNulls.java +++ b/java/core/src/test/org/apache/orc/TestRowFilteringComplexTypesNulls.java @@ -170,7 +170,8 @@ public void readEverything() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p >= 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p >= 0.03); } @Test @@ -264,7 +265,8 @@ public void readEverythingWithFilter() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p >= 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p >= 0.03); } @Test @@ -329,7 +331,8 @@ public void filterWithSeek() throws IOException { } FileSystem.Statistics stats = readEnd(); double readPercentage = readPercentage(stats, fs.getFileStatus(filePath).getLen()); - assertTrue(readPercentage > 0.07); + // Threshold lowered from 0.07: footer dedup removes one footer read per stripe. + assertTrue(readPercentage > 0.03); } private void seekToRow(RecordReader rr, VectorizedRowBatch b, long row) throws IOException { diff --git a/java/core/src/test/org/apache/orc/TestRowFilteringIOSkip.java b/java/core/src/test/org/apache/orc/TestRowFilteringIOSkip.java index fd32a431d7..a984c1e0fc 100644 --- a/java/core/src/test/org/apache/orc/TestRowFilteringIOSkip.java +++ b/java/core/src/test/org/apache/orc/TestRowFilteringIOSkip.java @@ -262,7 +262,8 @@ public void readWithSArg() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p >= 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p >= 0.03); } @Test @@ -306,7 +307,8 @@ public void readWithInvalidSArgAs() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p > 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p > 0.03); } private long validateFilteredRecordReader(RecordReader rr, VectorizedRowBatch b) @@ -396,7 +398,8 @@ public void readEverything() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p >= 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p >= 0.03); } private double readPercentage(FileSystem.Statistics stats, long fileSize) { @@ -421,7 +424,8 @@ public void readEverythingWithFilter() throws IOException { } double p = readPercentage(readEnd(), fs.getFileStatus(filePath).getLen()); assertEquals(RowCount, rowCount); - assertTrue(p >= 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(p >= 0.03); } @Test @@ -438,7 +442,8 @@ public void filterAlternateBatches() throws IOException { } FileSystem.Statistics stats = readEnd(); double readPercentage = readPercentage(stats, fs.getFileStatus(filePath).getLen()); - assertTrue(readPercentage > 0.06); + // Threshold lowered from 0.06: footer dedup removes one footer read per stripe. + assertTrue(readPercentage > 0.03); assertTrue(RowCount > rowCount); } @@ -490,7 +495,8 @@ public void filterWithSeek() throws IOException { } FileSystem.Statistics stats = readEnd(); double readPercentage = readPercentage(stats, fs.getFileStatus(filePath).getLen()); - assertTrue(readPercentage > 0.07); + // Threshold lowered from 0.07: footer dedup removes one footer read per stripe. + assertTrue(readPercentage > 0.03); } @Test diff --git a/java/core/src/test/org/apache/orc/impl/TestRecordReaderImpl.java b/java/core/src/test/org/apache/orc/impl/TestRecordReaderImpl.java index 3c70b7284a..f90f6d54a0 100644 --- a/java/core/src/test/org/apache/orc/impl/TestRecordReaderImpl.java +++ b/java/core/src/test/org/apache/orc/impl/TestRecordReaderImpl.java @@ -49,12 +49,17 @@ import org.apache.orc.OrcProto; import org.apache.orc.Reader; import org.apache.orc.RecordReader; +import org.apache.orc.StripeInformation; import org.apache.orc.TestConf; import org.apache.orc.TestVectorOrcFile; import org.apache.orc.TypeDescription; import org.apache.orc.Writer; +import org.apache.orc.impl.DataReaderProperties; +import org.apache.orc.impl.InStream; +import org.apache.orc.impl.OrcCodecPool; import org.apache.orc.impl.RecordReaderImpl.Location; import org.apache.orc.impl.RecordReaderImpl.SargApplier; +import org.apache.orc.impl.RecordReaderUtils; import org.apache.orc.impl.reader.ReaderEncryption; import org.apache.orc.impl.reader.StripePlanner; import org.apache.orc.impl.reader.tree.TypeReader; @@ -2817,4 +2822,96 @@ public void testDecimalIsRepeatingFlag() throws IOException { } } } + + /** + * Regression test for the duplicate stripe-footer read. When scanning the + * stripes of an ORC file, {@link RecordReaderImpl#readStripe()} reads the + * footer via {@link RecordReaderImpl#readStripeFooter(StripeInformation)} in + * {@code beginReadStripe} and must reuse that already-read footer when it + * calls {@link StripePlanner#parseStripe} instead of reading it a second + * time. This asserts exactly one {@code readStripeFooter} call per stripe + * for both a plain scan and a SArg-filtered scan. + */ + @Test + public void testNoDuplicateStripeFooterRead() throws Exception { + Path path = new Path(workDir, "noDupFooterRead.orc"); + FileSystem fs = FileSystem.get(conf); + fs.delete(path, true); + + TypeDescription schema = TypeDescription.fromString("struct"); + // rowIndexStride of 10_000 produces one row group per stripe so that SArg + // evaluation still keeps every stripe rather than skipping whole stripes. + int rowIndexStride = 10_000; + int rowsPerStripe = 10_000; + int stripeCount = 3; + conf.setLong(OrcConf.STRIPE_ROW_COUNT.getAttribute(), rowsPerStripe); + OrcFile.WriterOptions writerOptions = + OrcFile.writerOptions(conf).setSchema(schema) + .rowIndexStride(rowIndexStride); + Writer writer = OrcFile.createWriter(path, writerOptions); + VectorizedRowBatch writeBatch = schema.createRowBatch(); + LongColumnVector writeX = (LongColumnVector) writeBatch.cols[0]; + int base = 0; + int remaining = stripeCount * rowsPerStripe; + while (remaining > 0) { + writeBatch.reset(); + int chunk = Math.min(writeBatch.getMaxSize(), remaining); + for (int r = 0; r < chunk; ++r) { + int idx = writeBatch.size++; + writeX.vector[idx] = ++base; + } + writer.addRowBatch(writeBatch); + remaining -= chunk; + } + writer.close(); + // Reset the conf so other tests are not affected. + conf.set(OrcConf.STRIPE_ROW_COUNT.getAttribute(), + String.valueOf(OrcConf.STRIPE_ROW_COUNT.getDefaultValue())); + + Reader reader = OrcFile.createReader(path, OrcFile.readerOptions(conf).filesystem(fs)); + assertEquals(stripeCount, reader.getStripes().size()); + + // Plain scan: no SArg. + runStripeFooterCountTest(reader, stripeCount, null); + + // SArg scan: a predicate that evaluates to YES_NO for every row group so + // pickRowGroups is exercised but no stripe is skipped. + SearchArgument sarg = SearchArgumentFactory.newBuilder() + .startAnd() + .lessThan("x", PredicateLeaf.Type.LONG, (long) (stripeCount * rowsPerStripe + 1)) + .end().build(); + Reader sargReader = OrcFile.createReader(path, OrcFile.readerOptions(conf).filesystem(fs)); + runStripeFooterCountTest(sargReader, stripeCount, sarg); + } + + private void runStripeFooterCountTest(Reader reader, int stripeCount, + SearchArgument sarg) throws IOException { + ReaderImpl impl = (ReaderImpl) reader; + InStream.StreamOptions unencryptedOptions = + InStream.options() + .withCodec(OrcCodecPool.getCodec(reader.getCompressionKind())) + .withBufferSize(reader.getCompressionSize()); + DataReaderProperties properties = + DataReaderProperties.builder() + .withCompression(unencryptedOptions) + .withFileSystemSupplier(impl.getFileSystemSupplier()) + .withPath(impl.path) + .withZeroCopy(false) + .build(); + DataReader spied = spy(RecordReaderUtils.createDefaultDataReader(properties)); + // RecordReaderImpl clones the supplied DataReader before use, so make the + // spy observe all readStripeFooter calls by returning itself from clone(). + when(spied.clone()).thenReturn(spied); + Reader.Options readerOptions = reader.options().dataReader(spied); + if (sarg != null) { + readerOptions.searchArgument(sarg, new String[] {"x"}); + } + VectorizedRowBatch batch = reader.getSchema().createRowBatch(); + try (RecordReader rows = reader.rows(readerOptions)) { + while (rows.nextBatch(batch)) { + // drain + } + } + verify(spied, times(stripeCount)).readStripeFooter(any()); + } }