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);
+ }
}