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 af7f50e62e0 Fix stale broken cluster connection in CassandraIO (#39788)
af7f50e62e0 is described below

commit af7f50e62e05d743d31fce011ac2572ad3b67e4a
Author: Sharan Teja M <[email protected]>
AuthorDate: Tue Aug 18 14:34:42 2026 +0530

    Fix stale broken cluster connection in CassandraIO (#39788)
    
    * Fix for CassandraIO read  connection issue
    
    * Fixed a presubmit failure
    
    * Added test case
---
 .../beam/sdk/io/cassandra/ConnectionManager.java   | 15 ++++++++++-
 .../beam/sdk/io/cassandra/CassandraIOTest.java     | 30 ++++++++++++++++++++++
 2 files changed, 44 insertions(+), 1 deletion(-)

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..8a60275c707 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
@@ -59,9 +59,22 @@ public class ConnectionManager {
   }
 
   static 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