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

damccorm pushed a commit to branch release-2.76
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/release-2.76 by this push:
     new 161ea0b7b49 Fix stale broken cluster connection in CassandraIO 
(#39788) (#39803)
161ea0b7b49 is described below

commit 161ea0b7b49cd1894f9d293687955f3f82d20d07
Author: Sharan Teja M <[email protected]>
AuthorDate: Thu Aug 20 19:39:25 2026 +0530

    Fix stale broken cluster connection in CassandraIO (#39788) (#39803)
    
    * Fix stale broken cluster connection in CassandraIO (#39788)
    
    * Fix for CassandraIO read  connection issue
    
    * Fixed a presubmit failure
    
    * Added test case
    
    * Set getSession to synchronized to avoid race conditions
---
 .../beam/sdk/io/cassandra/ConnectionManager.java   | 17 ++++++++++--
 .../beam/sdk/io/cassandra/CassandraIOTest.java     | 30 ++++++++++++++++++++++
 2 files changed, 45 insertions(+), 2 deletions(-)

diff --git 
a/sdks/java/io/cassandra/src/main/java/org/apache/beam/sdk/io/cassandra/ConnectionManager.java
 
b/sdks/java/io/cassandra/src/main/java/org/apache/beam/sdk/io/cassandra/ConnectionManager.java
index 962e8ad8ec0..c2fb2f56d4e 100644
--- 
a/sdks/java/io/cassandra/src/main/java/org/apache/beam/sdk/io/cassandra/ConnectionManager.java
+++ 
b/sdks/java/io/cassandra/src/main/java/org/apache/beam/sdk/io/cassandra/ConnectionManager.java
@@ -58,10 +58,23 @@ public class ConnectionManager {
     return readToClusterHash(read) + read.keyspace().get();
   }
 
-  static Session getSession(Read<?> read) {
+  static synchronized Session getSession(Read<?> read) {
+    String clusterHash = readToClusterHash(read);
+    String sessionHash = readToSessionHash(read);
+
+    Cluster cachedCluster = clusterMap.get(clusterHash);
+
+    if (cachedCluster != null && cachedCluster.isClosed()) {
+      Session brokenSession = sessionMap.get(sessionHash);
+      if (brokenSession != null) {
+        sessionMap.remove(sessionHash, brokenSession);
+      }
+      // Removing broken cluster object
+      clusterMap.remove(clusterHash, cachedCluster);
+    }
     Cluster cluster =
         clusterMap.computeIfAbsent(
-            readToClusterHash(read),
+            clusterHash,
             k ->
                 CassandraIO.getCluster(
                     Objects.requireNonNull(read.hosts()),
diff --git 
a/sdks/java/io/cassandra/src/test/java/org/apache/beam/sdk/io/cassandra/CassandraIOTest.java
 
b/sdks/java/io/cassandra/src/test/java/org/apache/beam/sdk/io/cassandra/CassandraIOTest.java
index f63c819d420..93ba98af674 100644
--- 
a/sdks/java/io/cassandra/src/test/java/org/apache/beam/sdk/io/cassandra/CassandraIOTest.java
+++ 
b/sdks/java/io/cassandra/src/test/java/org/apache/beam/sdk/io/cassandra/CassandraIOTest.java
@@ -19,6 +19,9 @@ package org.apache.beam.sdk.io.cassandra;
 
 import static junit.framework.TestCase.assertTrue;
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNotSame;
 import static org.junit.Assert.assertNull;
 
 import com.datastax.driver.core.Cluster;
@@ -1218,4 +1221,31 @@ public class CassandraIOTest implements Serializable {
       return Objects.hashCode(tableColumn, indexColumn, valueColumn, data);
     }
   }
+
+  @Test
+  public void testSessionEvictionOnClosedCluster() {
+    CassandraIO.Read<String> readConfig =
+        CassandraIO.<String>read()
+            .withHosts(Collections.singletonList(CASSANDRA_HOST))
+            .withPort(cassandraPort)
+            .withKeyspace(CASSANDRA_KEYSPACE)
+            .withTable(CASSANDRA_TABLE);
+
+    Session initialSession = ConnectionManager.getSession(readConfig);
+    Cluster initialCluster = initialSession.getCluster();
+
+    initialCluster.close();
+    assertTrue("Cluster should be closed", initialCluster.isClosed());
+
+    Session newSession = ConnectionManager.getSession(readConfig);
+    Cluster newCluster = newSession.getCluster();
+
+    assertNotNull("New session should not be null", newSession);
+    assertFalse("New cluster should be open", newCluster.isClosed());
+
+    assertNotSame(
+        "ConnectionManager should create a new Session instance", 
initialSession, newSession);
+    assertNotSame(
+        "ConnectionManager should create a new Cluster instance", 
initialCluster, newCluster);
+  }
 }

Reply via email to