Skip to content

feat(core): detect corrupt log blocks and resume at the next one - #667

Open
linliu-code wants to merge 5 commits into
apache:mainfrom
linliu-code:feat/log-file-corruption-recovery
Open

feat(core): detect corrupt log blocks and resume at the next one#667
linliu-code wants to merge 5 commits into
apache:mainfrom
linliu-code:feat/log-file-corruption-recovery

Conversation

@linliu-code

Copy link
Copy Markdown

Stacked on #639#665 — review only the last commit.

What was there

fn create_corrupted_block_if_needed(&mut self, _pos: u64, _len: Option<u64>) -> Option<LogBlock> {
    // TODO: support creating corrupted block
    None
}

A corrupt or truncated block failed the entire log file. BlockType::Corrupted existed in the model but nothing could ever produce one.

What it does now

A well-formed Hudi log block records its total size twice — in the header length field and again in a trailing reverse pointer — and is followed by either another block or EOF. Three checks:

  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_add/checked_sub, so a garbage length reports corruption rather than panicking or allocating against it. That is why the check runs before the body is parsed rather than after — a bogus content length would otherwise reach a vec![0u8; len].

When a block fails, the reader emits a corrupt marker and resumes at the next magic marker, found by scanning in 1 MB windows that overlap by MAGIC.len() - 1 so a marker straddling a boundary is not missed. One bad block now costs its own span instead of the rest of the file.

Adopted from the internal reader, which has carried this and the recovery tests for some time.

Tests

Two harness cases un-ignored and passing — a corrupt tail block, and a delete-ordering fixture that was failing behind the same read.

Three unit tests on the check itself, since the harness cases only exercise one shape:

  • a length disagreeing with the trailing pointer is corrupt, and the real length is not — the negative case matters, since a check that always says "corrupt" would pass the positive one
  • a length past EOF, and u64::MAX, are corrupt — decided arithmetically, never by reading there
  • the recovery offset always lands within the file

Full workspace green: 1185 lib + 79 table-read + 39 datafusion + 21 + 12. Ignored 7 → 5.

Note

While verifying this I hit 117 unrelated failures that turned out to be a stale fixture-extraction cache under $TMPDIR/hudi-rs-test-fixtures, not the change — they reproduce on a clean tree. Clearing the directory fixes it. The cache key includes the zip's mtime, so it should self-invalidate; worth a look separately if it recurs in CI.

🤖 Generated with Claude Code

linliu-code and others added 5 commits August 3, 2026 22:01
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant