hudi-agent commented on code in PR #18776:
URL: https://github.com/apache/hudi/pull/18776#discussion_r4059139627
##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java:
##########
@@ -120,18 +121,31 @@ public void write() {
log.info("Schema evolution enabled for binary copy: {}",
schemaEvolutionEnabled);
records = this.writer.binaryCopy(inputFiles,
Collections.singletonList(path), writeScheMessageType, schemaEvolutionEnabled);
} catch (IOException e) {
+ closeWriterQuietly(e);
throw new HoodieIOException(e.getMessage(), e);
+ } catch (RuntimeException e) {
+ closeWriterQuietly(e);
+ throw e;
} finally {
this.recordsWritten = records;
this.insertRecordsWritten = records;
}
log.info("Finish rewriting {}. Using {} mills", this.path,
timer.endTimer());
}
+ private void closeWriterQuietly(Throwable failure) {
+ markClosed();
+ CloseableUtils.closeSuppressing(writer::close, failure);
Review Comment:
🤖 This now drives `HoodieParquetFileBinaryCopier.close()` on the failure
path, where `super.close()` → `writer.end()` runs *before*
`prefetchExecutor.shutdownNow()`. If `end()` throws on the half-written file
(fairly likely after a mid-row-group failure), the non-daemon single-thread
executor and any in-flight prefetch stream are never shut down. Could the
copier's `close()` shut down the executor in a `finally` (or before
`super.close()`) so this cleanup path doesn't leak the thread?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java:
##########
@@ -532,13 +532,14 @@ public Pair<ByteArrayOutputStream, Object>
serializeRecordsToLogBlock(HoodieStor
HoodieFileWriter parquetWriter = HoodieFileWriterFactory.getFileWriter(
HoodieFileFormat.PARQUET, outputStream, storage, config, writerSchema,
recordType);
- while (recordItr.hasNext()) {
- HoodieRecord record = recordItr.next();
- String recordKey = record.getRecordKey(readerSchema, keyFieldName);
- parquetWriter.write(recordKey, record, writerSchema);
+ try (HoodieFileWriter writerToClose = parquetWriter) {
+ while (recordItr.hasNext()) {
+ HoodieRecord record = recordItr.next();
Review Comment:
🤖 nit: `writerToClose` and `parquetWriter` refer to the same object inside
the try block — using two names for one variable is a bit confusing. Consider
just declaring `try (HoodieFileWriter parquetWriter =
HoodieFileWriterFactory.getFileWriter(...))` directly instead of introducing
the extra alias.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java:
##########
@@ -506,10 +502,32 @@ public List<WriteStatus> close() {
return Collections.singletonList(writeStatus);
} catch (IOException e) {
+ closeFileWriterQuietly(e);
throw new HoodieUpsertException("Failed to close UpdateHandle", e);
+ } catch (RuntimeException e) {
+ closeFileWriterQuietly(e);
+ throw e;
+ } finally {
+ keyToNewRecords = null;
Review Comment:
🤖 If `writeIncomingRecords()` or `closeFileWriter()` throws, the
`((Closeable) keyToNewRecords).close()` line above is skipped, and this
`finally` then drops the reference — so an `ExternalSpillableMap`'s on-disk
spill files are never cleaned up. Since `closeFileWriterQuietly` already
handles the file writer on failure, could it also `closeSuppressing`
`keyToNewRecords` when it's `Closeable`?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
##########
hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroHFileWriter.java:
##########
@@ -106,11 +107,17 @@ public HoodieAvroHFileWriter(String instantTime,
StoragePath file, HoodieHFileCo
.build();
StorageConfiguration<Configuration> storageConf = new
HadoopStorageConfiguration(conf);
StoragePath filePath = new StoragePath(this.file.toUri());
- OutputStream outputStream = HoodieStorageUtils.getStorage(filePath,
storageConf).create(filePath);
- this.writer = new HFileWriterImpl(context, outputStream);
- this.prevRecordKey = "";
- writer.appendFileInfo(
- HoodieAvroHFileReaderImplBase.SCHEMA_KEY,
getUTF8Bytes(schema.toString()));
+ OutputStream outputStream = HoodieStorageUtils.getStorage(filePath,
storageConf).create(filePath);
+ try {
+ this.writer = new HFileWriterImpl(context, outputStream);
+ this.prevRecordKey = "";
+ writer.appendFileInfo(
+ HoodieAvroHFileReaderImplBase.SCHEMA_KEY,
getUTF8Bytes(schema.toString()));
+ } catch (RuntimeException e) {
Review Comment:
🤖 nit: the ternary `writer != null ? writer : outputStream` to decide what
to close is a little subtle — worth a one-line comment explaining that `writer`
may not have been assigned yet if the `HFileWriterImpl` constructor itself
threw, so we fall back to closing the raw stream.
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]