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]

Reply via email to