This is an automated email from the ASF dual-hosted git repository.

stankiewicz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 6dcf56a4e96 Release HBaseIO resources even when an earlier close() 
throws (#39712)
6dcf56a4e96 is described below

commit 6dcf56a4e96859bce367c4979e15728d794b8048
Author: ZIHAN DAI <[email protected]>
AuthorDate: Mon Aug 31 19:06:01 2026 +1000

    Release HBaseIO resources even when an earlier close() throws (#39712)
    
    HBaseIO released its resources as consecutive, unguarded statements in
    three teardowns, so a throwing earlier close() skipped everything after
    it:
    
      HBaseReader.close()          scanner.close() -> connection.close()
      HBaseWriterFn.tearDown()     mutator.close() -> connection.close()
      WriteRowMutationsFn.tearDown()  table.close() -> 
HBaseSharedConnection.close()
    
    The writer case is not a corner case: BufferedMutator.close() is
    documented to perform a flush, so any bundle whose final buffered write
    fails already leaked the Connection it created in @Setup.
    
    The row-mutation case is worse than an ordinary leak. The skipped call
    is a reference-count decrement -- HBaseSharedConnection keeps a static
    pool, increments on getOrCreate and only closes the underlying
    Connection once the count reaches zero. Missing the decrement strands
    that entry, and its ZooKeeper session, for the lifetime of the JVM, and
    every later getOrCreate hands back the same unreleasable connection.
    
    Each teardown now runs every step and keeps the first failure, attaching
    later ones as suppressed. A plain nested try/finally would guarantee the
    calls happen but would silently swap which exception the caller sees, so
    it is only half a fix. The collection and rethrow logic is shared by the
    three sites as two package-private helpers on HBaseIO.
    
    HBaseReader, HBaseWriterFn and WriteRowMutationsFn drop `private` so the
    new test can construct them; they stay nested and are not part of any
    public API. Reverting each teardown individually fails exactly and only
    its own test.
    
    Fixes #39710
---
 sdks/java/io/hbase/build.gradle                    |   1 +
 .../java/org/apache/beam/sdk/io/hbase/HBaseIO.java |  91 ++++++++++--
 .../apache/beam/sdk/io/hbase/HBaseIOCloseTest.java | 160 +++++++++++++++++++++
 3 files changed, 241 insertions(+), 11 deletions(-)

diff --git a/sdks/java/io/hbase/build.gradle b/sdks/java/io/hbase/build.gradle
index c41e7edb317..7a06f6cc3cd 100644
--- a/sdks/java/io/hbase/build.gradle
+++ b/sdks/java/io/hbase/build.gradle
@@ -46,6 +46,7 @@ dependencies {
   testImplementation project(path: ":sdks:java:core", configuration: 
"shadowTest")
   testImplementation library.java.junit
   testImplementation library.java.hamcrest
+  testImplementation library.java.mockito_core
   // shaded-testing-utils has shaded all Hadoop/HBase dependencies
   
testImplementation("org.apache.hbase:hbase-shaded-testing-util:$hbase_version")
   testRuntimeOnly project(path: ":runners:direct-java", configuration: 
"shadow")
diff --git 
a/sdks/java/io/hbase/src/main/java/org/apache/beam/sdk/io/hbase/HBaseIO.java 
b/sdks/java/io/hbase/src/main/java/org/apache/beam/sdk/io/hbase/HBaseIO.java
index bc575b50af5..b1efb411481 100644
--- a/sdks/java/io/hbase/src/main/java/org/apache/beam/sdk/io/hbase/HBaseIO.java
+++ b/sdks/java/io/hbase/src/main/java/org/apache/beam/sdk/io/hbase/HBaseIO.java
@@ -48,6 +48,7 @@ import org.apache.beam.sdk.values.KV;
 import org.apache.beam.sdk.values.PBegin;
 import org.apache.beam.sdk.values.PCollection;
 import org.apache.beam.sdk.values.PDone;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.hbase.HRegionLocation;
 import org.apache.hadoop.hbase.TableName;
@@ -187,6 +188,33 @@ public class HBaseIO {
   /** Disallow construction of utility class. */
   private HBaseIO() {}
 
+  /**
+   * Collects a teardown failure. The first one is the one the caller ends up 
seeing; later ones are
+   * attached to it as suppressed, so nothing is lost and nothing replaces the 
original.
+   */
+  @VisibleForTesting
+  static Throwable appendSuppressed(@Nullable Throwable existingFailure, 
Throwable newFailure) {
+    if (existingFailure == null) {
+      return newFailure;
+    }
+    existingFailure.addSuppressed(newFailure);
+    return existingFailure;
+  }
+
+  /**
+   * Rethrows a failure collected by {@link #appendSuppressed}, preserving its 
type where it can.
+   */
+  @VisibleForTesting
+  static void rethrowCloseFailure(Throwable failure) throws IOException {
+    if (failure instanceof IOException) {
+      throw (IOException) failure;
+    }
+    if (failure instanceof RuntimeException) {
+      throw (RuntimeException) failure;
+    }
+    throw new IOException(failure);
+  }
+
   /**
    * Creates an uninitialized {@link HBaseIO.Read}. Before use, the {@code 
Read} must be initialized
    * with a {@link HBaseIO.Read#withConfiguration(Configuration)} that 
specifies the HBase instance,
@@ -492,7 +520,7 @@ public class HBaseIO {
     }
   }
 
-  private static class HBaseReader extends BoundedSource.BoundedReader<Result> 
{
+  static class HBaseReader extends BoundedSource.BoundedReader<Result> {
     private HBaseSource source;
     private Connection connection;
     private ResultScanner scanner;
@@ -549,14 +577,28 @@ public class HBaseIO {
     @Override
     public void close() throws IOException {
       LOG.debug("Closing reader after reading {} records.", recordsReturned);
+      // Release everything even if an earlier step throws, and keep the first 
failure: a
+      // connection left behind here outlives the reader.
+      Throwable failure = null;
       if (scanner != null) {
-        scanner.close();
+        try {
+          scanner.close();
+        } catch (Exception e) {
+          failure = appendSuppressed(failure, e);
+        }
         scanner = null;
       }
       if (connection != null) {
-        connection.close();
+        try {
+          connection.close();
+        } catch (Exception e) {
+          failure = appendSuppressed(failure, e);
+        }
         connection = null;
       }
+      if (failure != null) {
+        rethrowCloseFailure(failure);
+      }
     }
 
     @Override
@@ -732,7 +774,7 @@ public class HBaseIO {
 
     private final String tableId;
 
-    private class HBaseWriterFn extends DoFn<Mutation, Void> {
+    class HBaseWriterFn extends DoFn<Mutation, Void> {
 
       HBaseWriterFn(Write write) {
         checkNotNull(write.tableId, "tableId");
@@ -765,14 +807,28 @@ public class HBaseIO {
 
       @Teardown
       public void tearDown() throws Exception {
+        // BufferedMutator.close() performs a flush, so a failed final batch 
is an expected way
+        // for this to throw. Release the connection anyway, and keep the 
flush failure.
+        Throwable failure = null;
         if (mutator != null) {
-          mutator.close();
+          try {
+            mutator.close();
+          } catch (Exception e) {
+            failure = appendSuppressed(failure, e);
+          }
           mutator = null;
         }
         if (connection != null) {
-          connection.close();
+          try {
+            connection.close();
+          } catch (Exception e) {
+            failure = appendSuppressed(failure, e);
+          }
           connection = null;
         }
+        if (failure != null) {
+          rethrowCloseFailure(failure);
+        }
       }
 
       @Override
@@ -900,7 +956,7 @@ public class HBaseIO {
     private final String tableId;
 
     /** Function to write row mutations to a hbase table. */
-    private class WriteRowMutationsFn extends DoFn<KV<byte[], RowMutations>, 
Integer> {
+    class WriteRowMutationsFn extends DoFn<KV<byte[], RowMutations>, Integer> {
 
       public WriteRowMutationsFn(WriteRowMutations writeRowMutations) {
         checkNotNull(writeRowMutations.tableId, "tableId");
@@ -930,13 +986,26 @@ public class HBaseIO {
 
       @Teardown
       public void tearDown() throws Exception {
-
+        // HBaseSharedConnection.close() is a reference-count decrement, not 
an ordinary close.
+        // Skipping it strands the entry in the static pool for the lifetime 
of the JVM, so it has
+        // to run even when the table fails to close.
+        Throwable failure = null;
         if (table != null) {
-          table.close();
+          try {
+            table.close();
+          } catch (Exception e) {
+            failure = appendSuppressed(failure, e);
+          }
           table = null;
         }
-
-        HBaseSharedConnection.close(configuration);
+        try {
+          HBaseSharedConnection.close(configuration);
+        } catch (Exception e) {
+          failure = appendSuppressed(failure, e);
+        }
+        if (failure != null) {
+          rethrowCloseFailure(failure);
+        }
       }
 
       @ProcessElement
diff --git 
a/sdks/java/io/hbase/src/test/java/org/apache/beam/sdk/io/hbase/HBaseIOCloseTest.java
 
b/sdks/java/io/hbase/src/test/java/org/apache/beam/sdk/io/hbase/HBaseIOCloseTest.java
new file mode 100644
index 00000000000..d4aee9a7c7d
--- /dev/null
+++ 
b/sdks/java/io/hbase/src/test/java/org/apache/beam/sdk/io/hbase/HBaseIOCloseTest.java
@@ -0,0 +1,160 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.hbase;
+
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertThrows;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+import java.io.IOException;
+import java.lang.reflect.Field;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.hbase.HBaseConfiguration;
+import org.apache.hadoop.hbase.client.BufferedMutator;
+import org.apache.hadoop.hbase.client.Connection;
+import org.apache.hadoop.hbase.client.ResultScanner;
+import org.apache.hadoop.hbase.client.Table;
+import org.junit.After;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * {@link HBaseIO} used to release its resources as consecutive, unguarded 
statements, so a throwing
+ * earlier {@code close()} skipped everything after it.
+ *
+ * <p>The sharpest case is {@link 
HBaseIO.WriteRowMutations.WriteRowMutationsFn#tearDown()}: the
+ * skipped call there is {@link HBaseSharedConnection#close(Configuration)}, 
which is a
+ * reference-count decrement rather than an ordinary close. Missing it strands 
the entry in a {@code
+ * static} pool for the lifetime of the JVM, so that test asserts the real 
count rather than a mock
+ * interaction.
+ */
+@RunWith(JUnit4.class)
+public class HBaseIOCloseTest {
+
+  private final Configuration configuration = HBaseConfiguration.create();
+
+  @After
+  public void resetConnectionPool() throws IOException {
+    HBaseSharedConnection.closeAll();
+  }
+
+  // ---------------------------------------------------------------- failure 
collection
+
+  @Test
+  public void appendSuppressedKeepsTheFirstFailure() {
+    IOException first = new IOException("first");
+    IOException second = new IOException("second");
+
+    assertSame(first, HBaseIO.appendSuppressed(null, first));
+
+    Throwable collected = HBaseIO.appendSuppressed(first, second);
+    assertSame(first, collected);
+    assertArrayEquals(new Throwable[] {second}, collected.getSuppressed());
+  }
+
+  @Test
+  public void rethrowCloseFailurePreservesTheType() {
+    IOException io = new IOException("io");
+    assertSame(io, assertThrows(IOException.class, () -> 
HBaseIO.rethrowCloseFailure(io)));
+
+    IllegalStateException unchecked = new IllegalStateException("unchecked");
+    assertSame(
+        unchecked,
+        assertThrows(IllegalStateException.class, () -> 
HBaseIO.rethrowCloseFailure(unchecked)));
+
+    // Anything else has to be wrapped, because the teardowns only declare 
IOException.
+    Throwable checked = new Exception("checked");
+    IOException wrapped =
+        assertThrows(IOException.class, () -> 
HBaseIO.rethrowCloseFailure(checked));
+    assertSame(checked, wrapped.getCause());
+  }
+
+  // ---------------------------------------------------------------- reader
+
+  @Test
+  public void readerClosesTheConnectionWhenTheScannerFailsToClose() throws 
Exception {
+    HBaseIO.Read read = 
HBaseIO.read().withConfiguration(configuration).withTableId("some_table");
+    HBaseIO.HBaseReader reader =
+        new HBaseIO.HBaseReader(new HBaseIO.HBaseSource(read, null /* 
estimatedSizeBytes */));
+
+    ResultScanner scanner = mock(ResultScanner.class);
+    Connection connection = mock(Connection.class);
+    IOException scannerFailure = new IOException("scanner close failed");
+    doThrow(scannerFailure).when(scanner).close();
+    // The reader only acquires these in start(), which would need a live 
cluster.
+    set(reader, "scanner", scanner);
+    set(reader, "connection", connection);
+
+    assertSame(scannerFailure, assertThrows(IOException.class, reader::close));
+    verify(connection).close();
+  }
+
+  // ---------------------------------------------------------------- mutation 
writer
+
+  @Test
+  public void writerClosesTheConnectionWhenTheFinalFlushFails() throws 
Exception {
+    HBaseIO.Write write =
+        
HBaseIO.write().withConfiguration(configuration).withTableId("some_table");
+    HBaseIO.Write.HBaseWriterFn fn = write.new HBaseWriterFn(write);
+
+    BufferedMutator mutator = mock(BufferedMutator.class);
+    Connection connection = mock(Connection.class);
+    // BufferedMutator.close() performs a flush, so this is the expected 
failure mode.
+    IOException flushFailure = new IOException("flush on close failed");
+    doThrow(flushFailure).when(mutator).close();
+    set(fn, "mutator", mutator);
+    set(fn, "connection", connection);
+
+    assertSame(flushFailure, assertThrows(IOException.class, fn::tearDown));
+    verify(connection).close();
+  }
+
+  // ---------------------------------------------------------------- 
row-mutation writer
+
+  @Test
+  public void 
rowMutationWriterReleasesTheSharedConnectionWhenTheTableFailsToClose()
+      throws Exception {
+    HBaseSharedConnection.getOrCreate(configuration);
+    assertEquals(1, HBaseSharedConnection.getConnectionCount(configuration));
+
+    HBaseIO.WriteRowMutations write =
+        
HBaseIO.writeRowMutations().withConfiguration(configuration).withTableId("some_table");
+    HBaseIO.WriteRowMutations.WriteRowMutationsFn fn = write.new 
WriteRowMutationsFn(write);
+
+    Table table = mock(Table.class);
+    IOException tableFailure = new IOException("table close failed");
+    doThrow(tableFailure).when(table).close();
+    set(fn, "table", table);
+
+    assertSame(tableFailure, assertThrows(IOException.class, fn::tearDown));
+    // The point of the fix: the reference count still went back down, so the 
pooled connection is
+    // releasable instead of being stranded for the lifetime of the JVM.
+    assertEquals(0, HBaseSharedConnection.getConnectionCount(configuration));
+  }
+
+  private static void set(Object target, String fieldName, Object value) 
throws Exception {
+    Field field = target.getClass().getDeclaredField(fieldName);
+    field.setAccessible(true);
+    field.set(target, value);
+  }
+}

Reply via email to