Skip to content

GH-3706: Mark writer aborted on end() failure to avoid flushing incomplete files - #3707

Open
arg3t wants to merge 2 commits into
apache:masterfrom
arg3t:footer-corruption-on-write-fix
Open

arg3t wants to merge 2 commits into
apache:masterfrom
arg3t:footer-corruption-on-write-fix

Conversation

@arg3t

@arg3t arg3t commented Aug 4, 2026 •

Copy link
Copy Markdown

Rationale for this change

Makes ParquetFileWriter#end() mark the writer aborted before flushing, so an incomplete file is never committed on failure. Previously close() ran before the catch block in withAbortOnFailure and withAbortOnFailure only caught IOException, causing:

  • an incomplete file being flushed to storage when serializing indexes/stats/footer throws any exception
  • a RuntimeException (e.g. OOM) to be missed, leaving the writer un-aborted and still flushing

What changes are included in this PR?

  • withAbortOnFailure catches Throwable (not just IOException) and marks the writer aborted.
  • end() runs close() in an outer finally, so abort happens before close() and the flush is skipped on failure.

Are these changes tested?

Yes. 2 new tests in TestParquetFileWriter inject an IOException/RuntimeException during end() and assert the stream is never flushed.

Are there any user-facing changes?

No. Successful writes are unchanged; only the failure path is affected.

Closes #3706

@Fokko Fokko left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Makes sense, thanks @arg3t for fixing this!

@Fokko Fokko added this to the 1.19.0 milestone Sep 10, 2026
@jenwitteng

Copy link
Copy Markdown

Thanks, the end() reordering makes sense. There's a related gap on the InternalParquetRecordWriter side that this PR doesn't reach. It isn't blocking, and I can open a separate issue and PR if you'd rather keep this one as is.

InternalParquetRecordWriter.close() still catches only Exception:

} catch (Exception e) {
parquetFileWriter.abort();
throw e;
} finally {
AutoCloseables.uncheckedClose(columnStore, pageStore, bloomFilterWriteStore, parquetFileWriter);
closed = true;

An Error thrown in close() outside any ParquetFileWriter method never goes through withAbortOnFailure. So aborted stays false, and the finally calls ParquetFileWriter.close(), which flushes and closes out. This happens in two places:

  • flushRowGroupToStore() → columnStore.flush() (L231). The last row group's pages are encoded and compressed in ColumnChunkPageWriter.writePage, after startBlock() and before flushToFileWriter().
  • writeSupport.finalizeWrite() (L138).

On storage where the object only appears once the stream is closed (e.g. S3A), this publishes a footerless file, the same outcome as #3706.

I reproduced it on 1.17.0, which has #3351 but not this PR. The test uses an OutputFile whose stream close() marks the object as committed:

Failure inside close() 1.17.0 1.17.0 + this PR's withAbortOnFailure change 1.17.0 + catch (Throwable t) in IPRW
UnsatisfiedLinkError: snappy-java native load fails (no fault injection). The file is small, so the first compress() runs in close() committed, 4 bytes (PAR1 only) committed, 4 bytes not committed
OutOfMemoryError from the ByteBufferAllocator passed via withAllocator(), armed just before close() committed, 8832 bytes, no footer committed, 8832 bytes, no footer not committed
RuntimeException from the same allocator not committed not committed not committed

None of the stacks contains a ParquetFileWriter frame, and end() is never reached. For the first row the stack is SnappyCompressor.maxCompressedLength ← CodecFactory$HeapBytesCompressor.compress ← ColumnChunkPageWriter.writePage ← ColumnWriteStoreBase.flush ← InternalParquetRecordWriter.flushRowGroupToStore ← InternalParquetRecordWriter.close. That is why the withAbortOnFailure change doesn't affect them. (The middle column comes from that change patched into the 1.17.0 classes, not from a build of this branch.)

write() already uses catch (Throwable t) (L164), so the matching fix would be:

} catch (Throwable t) {
  parquetFileWriter.abort();
  throw t;
}

Precise rethrow keeps the throws IOException, InterruptedException clause, and abort() only sets a flag. A test for this needs the failure to happen outside ParquetFileWriter, for example an allocator via withAllocator() or a WriteSupport whose finalizeWrite() throws an Error. FaultInjectingOutputFile doesn't exercise it, because its faults fire inside ParquetFileWriter methods, which this PR already covers.

Repro (no fault injection)

Run with -Dorg.xerial.snappy.disable.bundled.libs=true to simulate snappy-java failing to load its native library:

import java.io.ByteArrayOutputStream;
import org.apache.parquet.example.data.Group;
import org.apache.parquet.example.data.simple.SimpleGroupFactory;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.example.ExampleParquetWriter;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.io.OutputFile;
import org.apache.parquet.io.PositionOutputStream;
import org.apache.parquet.schema.MessageType;
import org.apache.parquet.schema.MessageTypeParser;

public class SnappyLoadFailInClose {
  public static void main(String[] args) throws Exception {
    boolean[] committed = {false};
    ByteArrayOutputStream bytes = new ByteArrayOutputStream();
    OutputFile out = new OutputFile() {
      public PositionOutputStream create(long h) { return createOrOverwrite(h); }
      public PositionOutputStream createOrOverwrite(long h) {
        return new PositionOutputStream() {
          public long getPos() { return bytes.size(); }
          public void write(int b) { bytes.write(b); }
          public void write(byte[] b, int o, int l) { bytes.write(b, o, l); }
          public void close() { committed[0] = true; } // object stores commit here
        };
      }
      public boolean supportsBlockSize() { return false; }
      public long defaultBlockSize() { return 0; }
    };
    MessageType schema = MessageTypeParser.parseMessageType("message m { required int64 id; }");
    ParquetWriter<Group> w = ExampleParquetWriter.builder(out).withType(schema)
        .withCompressionCodec(CompressionCodecName.SNAPPY).build();
    SimpleGroupFactory f = new SimpleGroupFactory(schema);
    for (long i = 0; i < 100; i++) w.write(f.newGroup().append("id", i));
    try { w.close(); } catch (Throwable t) { System.out.println("close() threw " + t); }
    byte[] b = bytes.toByteArray();
    System.out.println("committed=" + committed[0] + " size=" + b.length + " footer="
        + (b.length >= 8 && new String(b, b.length - 4, 4, "US-ASCII").equals("PAR1")));
  }
}

Output on 1.17.0:

close() threw java.lang.UnsatisfiedLinkError: no snappyjava in java.library.path: ...
committed=true size=4 footer=false

arg3t pushed a commit to arg3t/parquet-java that referenced this pull request Oct 2, 2026
… to abort on Error

Address review feedback on apache#3707: close() caught only Exception, so an Error
thrown outside ParquetFileWriter (e.g. OOM during flushRowGroupToStore, or
finalizeWrite) left the writer un-aborted and the finally flushed an
incomplete, footerless file. Match write()'s catch (Throwable t) so abort
runs before close().
@arg3t

arg3t commented Oct 2, 2026

Copy link
Copy Markdown
Author

Nice catch @jenwitteng I made the adjustment you recommended. @Fokko any way you can take a look and approve?

… to abort on Error

Address review feedback on apache#3707: close() caught only Exception, so an Error
thrown outside ParquetFileWriter (e.g. OOM during flushRowGroupToStore, or
finalizeWrite) left the writer un-aborted and the finally flushed an
incomplete, footerless file. Match write()'s catch (Throwable t) so abort
runs before close().
@arg3t
arg3t force-pushed the footer-corruption-on-write-fix branch from 6454716 to d5af1ca Compare October 2, 2026 13:56
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.

ParquetFileWriter.end() flushes an incomplete file to storage when finalizing fails

3 participants