Conversation
|
Thanks, the
An
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
None of the stacks contains a
} catch (Throwable t) {
parquetFileWriter.abort();
throw t;
}Precise rethrow keeps the Repro (no fault injection)Run with 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: |
… 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().
|
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().
6454716 to
d5af1ca
Compare
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:
What changes are included in this PR?
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