HIVE-30059: LLAP IO Cache: native parquet data cache - #6793
abstractdog wants to merge 16 commits into
Conversation
| * Vectorized reads of parquet files from columns with list or map type is only supported if the nested types are of | ||
| * primitive type category | ||
| * check {@link VectorizedParquetRecordReader#checkListColumnSupport} for details on nested types under lists | ||
| * check {@link org.apache.hadoop.hive.ql.io.parquet.vector.ParquetRowGroupDecoder#checkListColumnSupport} for |
There was a problem hiding this comment.
do we need full name with package here?
| serdeCache, bufferManagerGeneric, conf, cacheMetrics, ioMetrics, tracePool, encodeExecutor) : null; | ||
| // Native Parquet IO is gated per query by the job conf at the dispatch sites. | ||
| this.parquetCvp = dataCache != null | ||
| ? new ParquetColumnVectorProducer(dataCache, bufferManagerOrc, conf, cacheMetrics, ioMetrics) |
There was a problem hiding this comment.
Should we drop orc from bufferManagerOrc ?
| private boolean checkOrcSchemaEvolution() { | ||
| SchemaEvolution evolution = rp.getSchemaEvolution(); | ||
| if (evolution == null) { | ||
| // No ORC-style schema evolution to validate (e.g. native parquet path); |
There was a problem hiding this comment.
why we check non-ORC schema inside checkOrcSchemaEvolution ?
|
|
||
| /** Slices of the cached buffers covering exactly the chunk's byte region, in file order. */ | ||
| private static List<ByteBuffer> chunkBuffers(ParquetEncodedColumnBatch batch, int pc) { | ||
| long start = batch.chunks[pc].getStartingPos(), end = start + batch.chunks[pc].getTotalSize(); |
There was a problem hiding this comment.
please move env var on a new line
| List<ByteBuffer> slices = new ArrayList<>(batch.columnBuffers[pc].length); | ||
| for (int i = 0; i < batch.columnBuffers[pc].length; ++i) { | ||
| long offset = batch.bufferOffsets[pc][i]; | ||
| long from = Math.max(start, offset), to = Math.min(end, offset + batch.bufferLengths[pc][i]); |
There was a problem hiding this comment.
same, new line for to, maybe need better naming
| super(consumer, includes.getPhysicalColumnIds().size(), ioMetrics, counters); | ||
| this.jobConf = jobConf; | ||
| this.useDecimal64ColumnVectors = HiveConf.getVar(jobConf, | ||
| ConfVars.HIVE_VECTORIZED_INPUT_FORMAT_SUPPORTS_ENABLED).equalsIgnoreCase("decimal_64"); |
There was a problem hiding this comment.
do we have constant for "decimal_64" ?
There was a problem hiding this comment.
🟡 Changes recommended
Unresolved critical and moderate findings affect schema projection, cache ownership, row lineage, configuration, and fallback behavior.
Get a fresh assessment by requesting another Copilot review.
Pull request overview
Adds opt-in native Parquet LLAP caching with vectored reads, cached-page decoding, fallback integration, and supporting tests.
Changes:
- Adds native Parquet LLAP configuration, APIs, and dispatch.
- Implements cache layouts, range buffers, vectored reads, and page decoding.
- Refactors shared Parquet reader logic and expands test coverage.
File summaries
| File | Reviewed changes and final comments |
|---|---|
ql/src/test/org/apache/hadoop/hive/llap/TestParquetRangeBuffers.java |
Tests range-buffer reuse and slicing; no final comments. |
ql/src/test/org/apache/hadoop/hive/llap/TestParquetCacheLayout.java |
Tests cache layout behavior; no final comments. |
ql/src/java/org/apache/hadoop/hive/ql/plan/MapWork.java |
Enables native Parquet LLAP eligibility; no final comments. |
ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/VectorizedParquetRecordReader.java |
Critical, 1 vote: Native reads omit row-lineage columns, producing null/default lineage values. |
ql/src/java/org/apache/hadoop/hive/ql/io/parquet/vector/ParquetRowGroupDecoder.java |
Shares row-group decoder construction; no final comments. |
ql/src/java/org/apache/hadoop/hive/ql/io/parquet/ParquetRecordReaderBase.java |
Reviewed as part of Parquet reader integration; no final comments. |
ql/src/java/org/apache/hadoop/hive/ql/io/HiveInputFormat.java |
Gates native Parquet wrapping; no final comments. |
ql/src/java/org/apache/hadoop/hive/llap/ParquetRangeBuffers.java |
Manages vectored-read buffers; no final comments. |
ql/src/java/org/apache/hadoop/hive/llap/ParquetCacheLayout.java |
Moderate, 1 vote: Native layout can create fallback-incompatible cache gaps and trigger an assertion instead of refetching. |
ql/src/java/org/apache/hadoop/hive/llap/LlapCacheAwareFs.java |
Supports cache-aware fallback reads; no final comments. |
llap-server/src/test/org/apache/hadoop/hive/llap/io/encoded/TestParquetEncodedDataReader.java |
Tests native reads and caching; no final comments. |
llap-server/src/test/org/apache/hadoop/hive/llap/io/api/impl/TestLlapRecordReader.java |
Tests LLAP reader behavior; no final comments. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/encoded/ParquetEncodedDataReader.java |
Critical, 2 votes: Top-level field indexes are used against flattened leaf columns, misselecting nested projections. Critical, 1 vote: Partial cache insertion failure can cause cache-owned buffers to be deallocated. Nit, 1 vote: count(*) fetches and caches unnecessary column chunks. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/encoded/ParquetEncodedColumnBatch.java |
Carries cached row-group data; no final comments. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/decode/ParquetEncodedDataConsumer.java |
Moderate, 1 vote: Decode exceptions are swallowed, allowing later I/O and resource use after failure. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/decode/ParquetColumnVectorProducer.java |
Creates native Parquet pipelines; no final comments. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/decode/ParquetCachedPageReadStore.java |
Reads cached Parquet pages; no final comments. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapRecordReader.java |
Supports producer fallback; no final comments. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapIoImpl.java |
Moderate, 1 vote: Discards supplied fileKey and tag. Moderate, 1 vote: Can create the native producer without a metadata cache in memory mode none. Moderate, 1 vote: Does not honor the per-query native Parquet disable flag. |
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapInputFormat.java |
Integrates LLAP input handling; no final comments. |
llap-client/src/java/org/apache/hadoop/hive/llap/io/api/LlapIo.java |
Extends the Parquet LLAP API; no final comments. |
iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/vector/HiveVectorizedReader.java |
Updates footer-cache API usage; no final comments. |
iceberg/iceberg-handler/src/main/java/org/apache/iceberg/mr/hive/HiveIcebergStorageHandler.java |
Reviewed for storage-handler integration; no final comments. |
common/src/java/org/apache/hadoop/hive/conf/HiveConf.java |
Adds native Parquet configuration; no final comments. |
Review details
Suppressed comments (6)
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapIoImpl.java:503
- This constructs a plain
FileSplitand silently discards thefileKeyandtagsupplied by the public API. The native reader then derives a different key from the job configuration/filesystem inloadFooter()and computes its own tag, so callers that provide a stable file identity cannot read or populate the intended data-cache entry or preserve the administrative cache tag. Propagate these arguments through the native reader, or remove them from the API instead of ignoring them.
FileSplit split = new FileSplit(path, offset, length, (String[]) null);
try {
LlapRecordReader rr = LlapRecordReader.create(conf, split, tableIncludedCols, HiveStringUtils.getHostname(),
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapIoImpl.java:285
dataCacheis non-null inLLAP_IO_MEMORY_MODE=nonebecauseLlapIoImplassigns aSimpleBufferManagerto it while leavingfileMetadataCachenull. With the native flag enabled, this condition still createsparquetCvp;ParquetEncodedDataReader.loadFooterthen callsgetParquetFooterBuffersFromCache, whose first operation requires a non-null metadata cache, so the LLAP input format fails instead of falling back to the normal Parquet reader. Gate this producer on the metadata cache as well (or make the native reader bypass the footer cache when it is unavailable).
this.parquetCvp = dataCache != null
? new ParquetColumnVectorProducer(dataCache, bufferManagerOrc, conf, cacheMetrics, ioMetrics)
: null;
llap-server/src/java/org/apache/hadoop/hive/llap/io/api/impl/LlapIoImpl.java:500
- The API contract says this method returns null when native Parquet IO is disabled, but
parquetCvpis constructed independently of the per-query flag and this check only tests for null. A caller can therefore sethive.llap.io.parquet.native.enabled=falseand still start the native reader, bypassing the configuration gate used byHiveInputFormat. IncludeConfVars.LLAP_IO_PARQUET_NATIVE_ENABLEDin this guard.
if (parquetCvp == null) {
return null;
}
llap-server/src/java/org/apache/hadoop/hive/llap/io/decode/ParquetEncodedDataConsumer.java:195
- These exceptions are reported but swallowed, so
EncodedDataConsumer.consumeDatareturns normally andParquetEncodedDataReader.performDataReadcontinues fetching and decoding later row groups before callingsetDone. A corrupt page in an early group can therefore keep doing I/O and holding cache resources after the downstream reader has failed. Propagate the decode failure after recording it (or letconsumeDataperform the error notification) so the reader's abort path stops the split.
} catch (IOException | RuntimeException e) {
// parquet-mr reports decode failures as runtime ParquetDecodingException.
LlapIoImpl.LOG.error("Parquet decodeBatch failed for rowGroup " + batch.rowGroupIx + " of " + path, e);
downstreamConsumer.setError(e);
} finally {
llap-server/src/java/org/apache/hadoop/hive/llap/io/encoded/ParquetEncodedDataReader.java:197
- For the
count(*)projection,colsToIncludecan be empty whilerequestedSchemastill contains the table fields; the shared decoder explicitly documents that all readers are null in this case. This code nevertheless derives every requested field as a projected leaf, so the native path fetches and caches every column chunk even though it only needs each row group's row count. Short-circuit the projected-leaf list whenincludes.getPhysicalColumnIds()is empty.
int[] projected = projectedLeaves(requestedSchema, fileSchema);
ql/src/java/org/apache/hadoop/hive/llap/ParquetCacheLayout.java:88
- This native layout is incompatible with the existing Parquet fallback cache writer.
LlapCacheAwareFs.getAndValidateMissingChunksrequires every missing gap to start and end onmaxAllocboundaries, while this method can split a 5 MiB chunk into 4 MiB + 1 MiB. If the 1 MiB native buffer is evicted and a later nested-type/fallback reader requests the chunk, the fallback sees a gap starting atchunkStart + 4 MiBand throws anAssertionErrorinstead of refetching it. Both paths need to share the same layout, or the fallback must accept arbitrary cache gaps.
public int[] bufferSizes(long length) {
int count = 0;
for (long left = length; left > 0; ++count) {
left -= left < minBuffer ? left : Math.min(maxBuffer, Long.highestOneBit(left));
}
- Files reviewed: 24/24 changed files
- Comments generated: 3
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| List<Integer> leaves = new ArrayList<>(); | ||
| for (Type field : requestedSchema.getFields()) { | ||
| if (fileSchema.containsField(field.getName())) { | ||
| leaves.add(fileSchema.getFieldIndex(field.getName())); | ||
| } | ||
| } | ||
| return leaves.stream().mapToInt(Integer::intValue).toArray(); |
There was a problem hiding this comment.
good catch, this would totally mess up columns in case of schema like:
message hive_schema {
optional group nested {
optional int32 a;
optional int32 b;
}
optional int32 x;
}
where fileSchema.getFieldIndex("x") = 1, but fileSchema.getColumns() → [nested.a, nested.b, x] (three leaves, indices 0, 1, 2), so block.getColumns().get(1) was used for x, but it's actually nested.b
| lowLevelCache.putFileData(fileKey, ranges, cached, 0, Priority.NORMAL, counters, cacheTag); | ||
| for (int i = 0; i < run.count; ++i) { | ||
| if (cached[i] != fresh[i]) { |
There was a problem hiding this comment.
ack, solved, and filed HIVE-30065 for the same problem in ORC
| columnReaders = new ParquetRowGroupDecoder(fileSchema, initialDefaults).buildColumnReaders( | ||
| pages, requestedSchema, columnTypesList, colsToInclude, | ||
| ColumnProjectionUtils.isReadAllColumns(jobConf), skipTimestampConversion, writerTimezone, | ||
| skipProlepticConversion, legacyConversionEnabled); |
There was a problem hiding this comment.
added a guard with good explanation regarding when it fires and about the consequences
…LLAP cache Whoever puts a Parquet column chunk in the LLAP cache decides how the file is cut into cacheable pieces. Today that is the cache-aware stream under the vectorized reader; the native reader added later does the same. If the two cut a file differently, a chunk cached by one is not reusable by the other and the same bytes end up held twice under different keys. Define the layout once, before either uses it: the piece sizes, the largest range a vectored read may ask for, and the pooled buffers those ranges are read into. Pieces are powers of two, largest first, so each fills its buddy allocation exactly. The cache-aware stream cached a chunk as one buffer and left the allocator to round it up, so a 5 Mb chunk took 8 Mb; over a mix of chunk sizes that lost about half the cache. Any contiguous run of pieces decomposes into itself, so a gap left by eviction re-caches on the same boundaries. The range bound is the larger of 8 Mb -- where reads against S3 stop getting faster, and where Trino also splits -- and the allocator's maximum, since a cache buffer is always read whole. The buffer pool is per stream: the pool never evicts, so one shared across the daemon would pin every executor's peak on the heap for the daemon's life.
… its own class VectorizedParquetRecordReader built its VectorizedColumnReader array inline: walk the requested schema, match each Hive type to a Parquet column, apply the column defaults a schema-evolved file needs, and decide which nested shapes are readable at all. That logic depends only on the schemas and the Hive types, not on how the pages were obtained, but it could not be called from anywhere else. Move it to ParquetRowGroupDecoder unchanged, so a reader that gets its pages from somewhere other than a file -- an LLAP cache-backed consumer, next -- builds the same readers from the same rules rather than growing a second copy that drifts. No behaviour change: the reader now delegates to the new class and the qtests are untouched. checkListColumnSupport moves with it, so the Iceberg storage handler's reference to it follows.
Parquet is cache-only in LLAP IO: the file bytes are cached, but every task re-runs the vectorized Parquet reader over them. ORC instead caches the column chunks it decoded and hands the consumer a column-vector batch. This adds the same three pieces for Parquet -- an encoded data reader, an encoded data consumer and a column vector producer -- so a Parquet scan is served from the cache in the same shape as ORC. The reader reads the footer, prunes row groups against the search argument and fetches only the projected column chunks, filling the cache with what it read. The consumer decodes a row group from a cache-backed InputFile, so a second reader of the same chunk decodes from memory and never touches the filesystem. Rather than a second set of readers, as ORC needed, the cached pages are served through parquet's own PageReadStore, so both paths run the same decode: the row-group logic moves out of VectorizedParquetRecordReader into ParquetRowGroupDecoder unchanged, and both callers use it. Off by default behind hive.llap.io.parquet.native.enabled. A projection that reaches into a nested type falls back to the vectorized Parquet reader, as does any error in the native path.
… summary (META_HIT, META_MISS) The summary's META_HIT / META_MISS columns are wired to METADATA_CACHE_HIT and METADATA_CACHE_MISS, which OrcEncodedDataReader bumps around its file tail lookups. The Parquet path uses the same FileMetadataCache (LlapIoImpl. getParquetFooterBuffersFromCache calls getFileMetadata / putFileMetadata on it) but never bumped the counters, so a Parquet-only run showed zeroes in those columns. Bump them from the same site. LlapIoImpl cannot see QueryFragmentCounters (it lives in llap-server; LlapIo in llap-client), so the impl reports the outcome through a nullable BooleanRef out-parameter -- the same signalling pattern storage-api already uses -- and ParquetEncodedDataReader.loadFooter increments the counter it already holds. The two non-native callers (VectorizedParquetRecordReader on the vectorized fallback, HiveVectorizedReader on the Iceberg path) do not have per-fragment counters plumbed and pass null; their footer traffic continues to go unaccounted, as it did before. TestParquetEncodedDataReader reads the file twice and asserts the second read reports one hit and zero misses; only the second read is asserted because the footer cache is a @BeforeClass singleton and any earlier test in the class may have populated it. Co-Authored-By: Claude Code <noreply@anthropic.com>
Verifies the descriptors and time counters wired into ParquetEncodedDataReader in the previous commit: SELECTED_ROWGROUPS, TOTAL_IO_TIME_NS, HDFS_TIME_NS, METADATA_CACHE_HIT/MISS, plus FILE and STRIPES descriptors visible via the QueryFragmentCounters summary string. Co-Authored-By: Claude Code <noreply@anthropic.com>
- HiveIcebergStorageHandler: import ParquetRowGroupDecoder so the javadoc
{@link} can use the short name.
- LlapIoImpl: rename local bufferManagerOrc -> bufferManagerData; the
buffer manager is now shared with the native Parquet producer, so the
Orc suffix is stale.
- LlapRecordReader: rename checkOrcSchemaEvolution -> checkSchemaEvolution
since the method already tolerates producers (like native Parquet) that
do not expose an ORC-style SchemaEvolution.
- ParquetCachedPageReadStore.chunkBuffers: split the two-per-line locals,
and rename from/to -> sliceStart/sliceEnd (with chunkStart/chunkEnd and
bufferStart/bufferEnd for the surrounding bounds).
Co-Authored-By: Claude Code <noreply@anthropic.com>
Three Critical findings on the native Parquet cache path:
- projectedLeaves used fileSchema.getFieldIndex to pick column chunks, but
BlockMetaData.getColumns() is in flat leaf order. A file schema like
'group nested {a,b}, x' projecting 'x' would hit nested.b instead of x.
Map each requested top-level primitive to its single-segment leaf in
fileSchema.getColumns() instead. Add a regression test.
- putColumn set part.owned only after processing all missing ranges for the
column, so a mid-run throw from putFileData (e.g. its length-mismatch
guard) left already-inserted cache buffers looking like raw allocations
and finishFetch would allocator.deallocate cache-owned memory. Insert
one range at a time and flip ownership per part.
- ParquetEncodedDataReader.loadFooter builds requestedSchema without the
row-lineage columns the fallback reader adds via
RowLineageUtils.getRequestedSchemaWithRowLineageColumns, so with row
lineage on the native path would silently emit nulls for
ROW__LINEAGE__ID / LAST__UPDATED__SEQUENCE__NUMBER. Detect that case in
the producer and fall back, matching the nested-projection fallback.
Co-Authored-By: Claude Code <noreply@anthropic.com>
Fix the mechanical issues flagged on files this PR modifies. Skips the 15 findings inside verbatim extractions from VectorizedParquetRecordReader (kept behavior-preserving) and the switch/case indentation reports. Notable non-trivial fixes: - LlapIo.llapVectorizedParquetReaderForPath: 9 params -> 3, collected into a new LlapParquetReadRequest record. - LlapInputFormat.getRecordReader: extract wrapOrFallback() to bring cognitive complexity under threshold. - ParquetEncodedDataReader: extract planMissRun() from planColumnChunk; drop unused 'includes' field/param, unused maxAlloc param, and Part.miss; narrow catch(Throwable) -> catch(Exception) at four sites and preserve InterruptedException on the thread instead of masking it. - ParquetCacheLayout / ParquetEncodedDataReader: replace loop-counter mutation with while loops. - CacheChunk: drop stale @VisibleForTesting; it is part of the cache's public read surface for both ORC and Parquet. - checkstyle/suppressions.xml: suppress VisibilityModifier for ParquetEncodedColumnBatch (matches parent EncodedColumnBatch pattern). - Delete stale checkListColumnSupport from VectorizedParquetRecordReader (moved to ParquetRowGroupDecoder). Style-only elsewhere: pattern-instanceof, split multi-decls, empty-body 'why' comments, header rewraps, paren-pad, restricted-identifier rename (record() -> recordRanges()). Verified: mvn checkstyle:check on llap-server (0 violations) and the existing test suites: TestParquetEncodedDataReader (27/27), TestLlapRecordReader (2/2), TestParquetRangeBuffers (2/2). Co-Authored-By: Claude Code <noreply@anthropic.com>
66f34dc to
65352ad
Compare
|
| LLAP_IO_ENCODE_ENABLED("hive.llap.io.encode.enabled", true, | ||
| "Whether LLAP should try to re-encode and cache data for non-ORC formats. This is used\n" + | ||
| "on LLAP Server side to determine if the infrastructure for that is initialized."), | ||
| LLAP_IO_PARQUET_NATIVE_ENABLED("hive.llap.io.parquet.native.enabled", false, |
There was a problem hiding this comment.
Wouldn't be this enabled by default?
| private boolean checkOrcSchemaEvolution() { | ||
| SchemaEvolution evolution = rp.getSchemaEvolution(); | ||
| if (evolution == null) { | ||
| // No ORC-style schema evolution to validate (e.g. native parquet path); |



What changes were proposed in this pull request?
Why are the changes needed?
Does this PR introduce any user-facing change?
How was this patch tested?