From 640417352b815ebbbd3bfc8ba351aac7388a62b0 Mon Sep 17 00:00:00 2001 From: sharantej Date: Mon, 17 Aug 2026 16:59:40 +0530 Subject: [PATCH 1/4] Fix for CassandraIO read connection issue --- .../beam/sdk/io/cassandra/ConnectionManager.java | 16 ++++++++++++++-- 1 file changed, 14 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 962e8ad8ec00..305bee1dad10 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,21 @@ private static String readToSessionHash(Read read) { } static Session getSession(Read read) { - Cluster cluster = + String clusterHash = readToClusterHash(read); + String sessionHash = readToSessionHash(read); + + Cluster cluster = clusterMap.get(clusterHash); + + if(cluster != null && cluster.isClosed()) { + Session brokenSession = sessionMap.get(sessionHash); + if (brokenSession != null) { + sessionMap.remove(sessionHash,brokenSession); + } + clusterMap.remove(clusterHash,cluster); + } + cluster = clusterMap.computeIfAbsent( - readToClusterHash(read), + clusterHash, k -> CassandraIO.getCluster( Objects.requireNonNull(read.hosts()), From df6dafc66812b3676e276cffb472d809324f8ac1 Mon Sep 17 00:00:00 2001 From: sharantej Date: Tue, 18 Aug 2026 05:34:04 +0530 Subject: [PATCH 2/4] Fixed a presubmit failure --- .../beam/sdk/io/cassandra/ConnectionManager.java | 11 ++++++----- 1 file changed, 6 insertions(+), 5 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 305bee1dad10..8a60275c707a 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 @@ -62,16 +62,17 @@ static Session getSession(Read read) { String clusterHash = readToClusterHash(read); String sessionHash = readToSessionHash(read); - Cluster cluster = clusterMap.get(clusterHash); + Cluster cachedCluster = clusterMap.get(clusterHash); - if(cluster != null && cluster.isClosed()) { + if (cachedCluster != null && cachedCluster.isClosed()) { Session brokenSession = sessionMap.get(sessionHash); if (brokenSession != null) { - sessionMap.remove(sessionHash,brokenSession); + sessionMap.remove(sessionHash, brokenSession); } - clusterMap.remove(clusterHash,cluster); + // Removing broken cluster object + clusterMap.remove(clusterHash, cachedCluster); } - cluster = + Cluster cluster = clusterMap.computeIfAbsent( clusterHash, k -> From 50a1515cac14bc1bc791e0ed121f43604109759d Mon Sep 17 00:00:00 2001 From: sharantej Date: Tue, 18 Aug 2026 06:18:58 +0530 Subject: [PATCH 3/4] Added test case --- .../sdk/io/cassandra/CassandraIOTest.java | 30 +++++++++++++++++++ 1 file changed, 30 insertions(+) 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 f63c819d4202..93ba98af6745 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 @@ 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 int hashCode() { return Objects.hashCode(tableColumn, indexColumn, valueColumn, data); } } + + @Test + public void testSessionEvictionOnClosedCluster() { + CassandraIO.Read readConfig = + CassandraIO.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); + } } From 54e03498e7c7eb5a876212e15e36fb63bc37e3ef Mon Sep 17 00:00:00 2001 From: sharantej Date: Wed, 19 Aug 2026 15:23:24 +0530 Subject: [PATCH 4/4] Set getSession to synchronized to avoid race conditions --- .../org/apache/beam/sdk/io/cassandra/ConnectionManager.java | 2 +- 1 file changed, 1 insertion(+), 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 8a60275c707a..c2fb2f56d4eb 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,7 +58,7 @@ private static String readToSessionHash(Read read) { 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);