fix(core): legacy 2-level lists, log-block predicates, and schema promotion - #669
Open
linliu-code wants to merge 11 commits into
Open
fix(core): legacy 2-level lists, log-block predicates, and schema promotion#669linliu-code wants to merge 11 commits into
linliu-code wants to merge 11 commits into
Conversation
Squashed view of apache#639-apache#662 for review. Not for merge — the reviewable increments are those PRs; this is the same code in one diff. Ports the merge-on-read file group reader from onehouseinc/hudi-rs-internal into hudi-core, wires it behind a switch that defaults to the reader that has always served reads, and brings its end-to-end test harness across. The ported reader is `pub(crate)` and reached only through `hoodie.read.merge.engine = v2`. Nothing changes for anyone who does not set it. It is not at parity yet: the gaps are pinned as ignored cases carrying their findings, and the outstanding decisions are in the PR descriptions. Four fixes land on paths the existing reader shares: decimals had no Arrow conversion, Avro timestamp logical types lost their UTC zone, the properties-escaped create schema was not being unescaped by one of its two consumers, and the base file reader had no way to accept a pushdown predicate. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A delete block whose ordering value is anything but a small integer fails to
read, with errors like `Union index 1490 out of bounds: 2`. The index is
nonsense because the byte stream is misaligned, not because a branch was chosen
wrongly.
`orderingVal` is declared here as a union of primitives. Hudi writes a union of
per-type wrapper records, and inserted `BooleanWrapper` at position 1, so every
position from `int` onward names a different type than this crate assumes.
Position 3 is `float` here and `LongWrapper` there. Reading a long as a float
consumes four bytes instead of two, and the next record's key length is read
from the middle of the previous value.
The delete block in `table_delete_ord_long` is exactly self-consistent under
Hudi's schema and not under this one:
04 | 02 02 34 | 02 00 | 06 c0 3e | 02 02 33 | 02 00 | 06 f0 2e | 00
two records, keys "4" and "3", ordering position 3, values 4000 and 3000
Decoded here, position 3 is a float, so `c0 3e 02 02` is eaten and `33` becomes
the next key's union index: zigzag 0x33 is -26, which is the reported error.
So this takes Hudi's schema. Two things follow from it:
The Arrow side wants a scalar, not a record with one field, so the wrapper is
unwrapped after the schema is narrowed — narrowing reads the position Hudi
wrote, and unwrapping rewrites it, so the order matters.
The wrapper is chosen for the value rather than for the column, so a table
whose ordering column is a long can still carry an `IntWrapper` for a small
value. The delete batch's ordering column is now cast to the type the data
schema declares. The old schema hid this by calling position 2 a long
regardless, which happened to match the two fixtures that exercise it.
`ArrayWrapper` orders by a list and is rejected in both places that read the
position, rather than mapped to something that would disagree with its value.
Older tables written with the primitive union are not supported. No fixture
here uses one, and the two shapes cannot be told apart from the bytes: where
they differ, decoding usually fails, and at position 2 both succeed and yield
the same number.
Un-ignores the eleven cases that pinned this.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
An Avro map was modelled as `Dictionary(Utf8, V)`. An Arrow dictionary key must be an integer, so that is not a valid type — and it does not reconcile against the `Map` a parquet base file carries, so any table with a map column fails to read once a log block has to be merged with its base file. Avro maps become `Map(key_value: struct<key: string, value: V>)`, matching what the parquet reader produces, so the two agree by name and by shape. The array side builds a `MapArray`. Entries are materialized as two-field records so the existing struct machinery builds both children, which is also why `child_schema_lookup` now registers those two positions — a struct-valued map needs its own fields resolvable underneath them. Entries are emitted in key order. Avro maps are unordered and the Arrow type says so, but a stable order keeps a read reproducible rather than dependent on hash iteration. This is on the shared conversion, so it fixes the existing read path too: a map column has never been readable through either reader. Un-ignores the case covering NULL elements inside containers. Two other cases that were pinned on this stay pinned, now on a decimal column reading as NULL — a separate gap this one was masking. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The merge-on-read reader was given its schema from `hoodie.table.create.schema`. That is the schema the table was created with, and Hudi treats it as a last resort: `TableSchemaResolver` reads the latest commit's metadata first, then a base file's footer, and only then falls back to the create schema. This crate's own `schema::resolver::resolve_data_schema` does the same. The reader was reaching past both for the weakest source. Reading the base file's own schema is what the existing path effectively does, so the two engines now start from the same types. It is also what the data actually has: under schema evolution the create schema is stale, and the engine evolves each batch to the required schema regardless. Three workarounds go with it. The create schema arrives as Java writes a properties file, with `:` escaped, so it had to be unescaped before it would parse as JSON. It carries no `_hoodie_*` columns, so those had to be prepended when the table populates them. And a table that never recorded one could not be read at all — which included every reader built from a bare base URI, the shape the cxx bridge uses. Slices with no log files now go through the engine too. They were held back because the create schema modelled a map as an invalid Arrow dictionary and every fixture here has a map column; the schema no longer comes from there, and the conversion itself is fixed separately. The engine reduces to a base file read, and the test asserting that the setting does not change such a read now compares the two engines rather than one path with itself. Reading the footer costs one request. The engine reads it again when it opens the file; collapsing the two is worth doing but is not this change. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Corruption detection was a stub returning `None`, so a corrupt or truncated block failed the whole log file. Adopted from the internal reader, which has carried this for a while. A well-formed block records its total size twice — once in the header length field and once in a trailing reverse pointer — and is followed by either another block or the end of the file. Three checks decide the question: 1. the trailing pointer has to lie inside the file 2. the size it records has to agree with the header 3. what follows the block has to be a magic marker or the end Every offset is computed with checked arithmetic, so a garbage length reports corruption rather than panicking or allocating against it. That is the point of doing this before parsing the body rather than after. A block that fails yields a corrupt marker, and the reader resumes at the next magic marker — found by scanning in windows that overlap by five bytes, so a marker straddling a window boundary is not missed. One bad block now costs its own span instead of the rest of the file. Un-ignores the two cases that pinned this. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A decimal column in an Avro log block read as entirely NULL. The vendored converter resolves a decimal from `Value::Bytes` or `Value::Fixed`, but a declared `logicalType: decimal` decodes to `Value::Decimal`, which falls through to `None` — and a `None` there is indistinguishable from a field that was genuinely null, so the column came back empty rather than failing. Rather than add the missing arm, this hands the work to `arrow-avro`, which already handles the logical types and the container types correctly. The map conversion fixed in the previous change came from the same converter, and the pattern of a silent wrong answer is what makes it worth replacing rather than patching. Two shapes had to be reconciled. A Hudi data block frames each datum with a four-byte length and no Avro framing, while `arrow-avro`'s decoder expects a Single Object Encoding prefix. The writer schema is registered once and the ten bytes that yields — marker plus fingerprint — are written ahead of each body, into a buffer reused across records. And `arrow-avro` spells a UTC timestamp's zone as the offset `+00:00` where parquet spells it `UTC`. The same zone, but Arrow compares timezones as strings, so a log batch would refuse to concatenate with the base batch it merges with. Decoded batches are relabelled to `UTC`; the values are already UTC instants, so nothing is converted. The decoder is also columnar. The path it replaces built an `apache_avro::Value` per record — allocating a `String` per string field and a `Vec` per collection — and then walked that to fill the columns. The vendored converter stays for now: delete blocks decode from values that are already parsed, which is a different shape. Un-ignores the two cases that pinned the decimal. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A base file with an `array<map>` column written by parquet-avro fails to read with `Map cannot be repeated`, and the file is unreadable by either reader. That writer defaults to `write-old-list-structure=true`, which encodes the column as a 2-level list whose element is a REPEATED map group. parquet-rs walks the LIST node, takes the repeated child for the list wrapper and the map for the element, and its `visit_map` rejects a REPEATED map unconditionally. Both arrow-rs and arrow-cpp reject the same physical schema, so the encoding really is legacy — but Hudi's own reader accepts it, and tables already written this way have to stay readable. The reject cannot be avoided at the Arrow level, so the parquet schema is rewritten before the Arrow build: the legacy repeated map becomes a synthetic repeated `list` wrapper around a required map. The footer parse never trips the reject — only the Arrow build does — so normalizing in between fixes every reader built from that metadata. Nothing about the data changes. Every leaf keeps its definition and repetition levels and its position in the DFS order, so the same bytes decode to the same values and the original row groups are reused as they are. The rewrite itself was already here with its own tests; it had no caller. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The same encoding the previous change handled for base files also appears in parquet log blocks, and fails the same way — `Map cannot be repeated` — because the block was built straight from its footer with no chance to rewrite the schema in between. Parse the footer, normalize the legacy list, then build the reader from that metadata, exactly as the base file path now does. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The base file read pushes a row filter down; a parquet log block read does not, so every row of every block is decoded and handed to the merge even when the predicate has already ruled the block out. The decision of whether pushing is sound stays with the merge-on-read reader, which is the only place that knows. A log record can update a row, so a predicate may only be evaluated before the merge when the merge cannot change its answer. Log blocks exist only on merge-on-read, so the copy-on-write arm of that gate is unreachable here and it reduces to a predicate over primary keys, which are immutable across upserts. Anything else is left for the post-merge filter, as before. The filter is carried as an optional builder rather than a filter, because the predicate has to be matched against each block's own schema, which does not exist until its footer is read. A builder that declines reads every row. Both the log file reader and the block decoder take it through a setter that defaults to unset, so the log scanner and every other caller are unaffected. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A partial-update block carries only the columns that were written, and the merge distinguishes "this column was not in the update" from "this column was set to null". The block's own schema is the only place that signal exists. Decoding is already against that schema and never against a wider table schema, so the narrowness survives — but nothing said so. Upstream needs an explicit check here because it decodes against a required schema; this crate does not, which makes the behavior correct by construction and therefore easy to lose. The test fails the moment a reader schema is supplied to the decoder without excluding blocks that carry `IsPartial`, which is what will happen when Avro resolution or extended promotion lands. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
A table whose column was promoted from int to long could not be read: the merge failed with `column 'num' type Int32 incompatible with target Int64`. Both the base file and the older log blocks predate the promotion, and neither was being brought up to the current type. Avro defines int to long as a promotion, so the log side is resolved while the block is read rather than reconciled afterwards. The required schema is handed to the decoder, which fills columns added since the block was written and delivers promoted columns in the promoted type. A block carrying `IsPartial` keeps decoding against its own schema, since the merge needs to know which columns the update actually set. The base file is parquet, so Avro resolution does not reach it. It goes through the batch evolution instead, which knew how to widen a float and not an integer. Widening an integer is exact, so it casts directly rather than through a string the way float to double has to for parity with Java. Un-ignores the promotion case, which reads the base parquet, a log block written before the promotion, and one written after, and compares the merged result against the Spark snapshot — including a value beyond what an int can hold. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stacked on #639–#668 — review only the last five commits.
Three changes to the parquet paths, each wiring up code that was already in the tree with tests but no caller.
1–2. Legacy 2-level lists (
fix)A parquet file with an
array<map>column written by parquet-avro fails to read:Both readers, both file kinds — base files and parquet log blocks. The file is simply unreadable.
That writer defaults to
write-old-list-structure=true, encoding the column as a 2-level legacy list whose element is a REPEATED map group:parquet-rs takes the repeated child for the list wrapper and the map for the element, then dispatches it to
visit_map, which rejects a REPEATED map unconditionally. arrow-rs and arrow-cpp both reject this schema, so the encoding genuinely is legacy — but Hudi's own reader accepts it, and tables already written this way must stay readable.The reject can't be avoided at the Arrow level, so the parquet schema is rewritten before the Arrow build: the legacy repeated map becomes a synthetic repeated
listwrapper around a required map. The footer parse never trips the reject; only the Arrow build does, so normalizing in between fixes every reader built from that metadata.Nothing about the data changes. Every leaf keeps its definition and repetition levels and its DFS position, so the same bytes decode to the same values and the original row groups are reused as they are.
Commit 1 is shared code — the legacy reader, the v2 engine, file listing, schema resolution and statistics all read base files through
ParquetBaseFileReader, so it fixes the reader in use today. Commit 2 is the same fix for parquet log blocks.3. Predicate pushdown into parquet log blocks (
feat)The base file read pushes a row filter down (#662); a parquet log block read did not, so every row of every block was decoded and handed to the merge even when the predicate had already ruled it out.
The safety decision stays with the merge-on-read reader, the only place that knows. A log record can update a row, so a predicate may be evaluated before the merge only when the merge cannot change its answer. Log blocks exist only on merge-on-read, so the copy-on-write arm of
can_push_row_filter()is unreachable here and the gate reduces entirely to a predicate over primary keys — immutable across upserts. Anything else is left for the post-merge filter.Carried as a builder rather than a filter, because the predicate has to be matched against each block's own schema, which doesn't exist until its footer is read. A builder that declines reads every row.
Both the log file reader and the block decoder take it through a setter defaulting to unset, so the log scanner and every other caller are unaffected.
4. A partial-update block stays narrow (
test)Such a block carries only the columns that were written, and the merge distinguishes "not in this update" from "set to null". The block's own schema is the only place that signal lives.
Decoding was already against that schema, so this adds no behavior — it pins it. Upstream needs an explicit check because it decodes against a required schema; this crate did not, which made the behavior correct by construction and therefore easy to lose. Commit 5 is exactly the change that would have lost it.
5. Resolve a log block up to the promoted schema (
fix)A table whose column was promoted
int→longcould not be read:Both the base file and the older log blocks predate the promotion, and neither was brought up to the current type. Two halves:
IsPartialare excluded, which is what commit 4 guards.floatbut not anint. Widening an integer is exact, so it casts directly rather than through a string the wayfloat→doublemust for parity with Java.The fixture reads a base parquet, a log block written before the promotion, and one written after, and compares the merged result against the Spark snapshot — including a value beyond what an
intcan hold.Tests
Two legacy-list reads against the real captured fixture — a genuine Hudi-written block in the rejected encoding — one per path.
One end-to-end pushdown case: the fixture holds keys
k1/k2/k3, and a_hoodie_record_key = k2predicate must return exactly that row. It exercises the whole chain — reader decides it's safe, hands it to the log file reader, which hands it to the block decoder — so it fails if any link drops it. That matters here: this is exactly how base-file pushdown was silently lost during the port.Three unit tests on the Avro resolution: int→long, the same through a
["null", T]union (the shape every real Hudi column has), and no reader schema leaving the writer's type alone.Both new behaviors verified non-vacuous by mutation:
Full workspace green: 1194 lib + 79 table-read + 39 datafusion + 21 + 12. No new clippy findings in the changed lines (the two in
log_record_reader.rsare at 218–258, pre-existing).🤖 Generated with Claude Code