From e4d22904ed78031ad95b667420ef8d1416b5cfaa Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Mon, 3 Aug 2026 18:04:00 +0800 Subject: [PATCH] fix(cassandra): infer local datacenter when unspecified Closes #976 --- .../dbx/agent/cassandra/CassandraAgent.java | 19 +++++++- .../agent/cassandra/CassandraAgentTest.java | 46 ++++++++++++++++++- 2 files changed, 62 insertions(+), 3 deletions(-) diff --git a/agents/drivers/cassandra/src/main/java/com/dbx/agent/cassandra/CassandraAgent.java b/agents/drivers/cassandra/src/main/java/com/dbx/agent/cassandra/CassandraAgent.java index 7d3aaf09e..f7cdc196f 100644 --- a/agents/drivers/cassandra/src/main/java/com/dbx/agent/cassandra/CassandraAgent.java +++ b/agents/drivers/cassandra/src/main/java/com/dbx/agent/cassandra/CassandraAgent.java @@ -20,6 +20,7 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; public final class CassandraAgent extends AbstractJdbcAgent { + private static final String DC_INFERRING_LOAD_BALANCING = "loadbalancing=DcInferringLoadBalancingPolicy"; private static final Pattern TARGET_PATTERN = Pattern.compile("target[\"']?\\s*[:=]\\s*[\"']?([\\w]+)"); @Override @@ -164,14 +165,30 @@ public final class CassandraAgent extends AbstractJdbcAgent { String keyspace = coalesce(params.getDatabase()).trim(); // Cassandra rejects an empty keyspace path; omit it so DBX can connect first and list keyspaces. String url = keyspace.isEmpty() ? baseUrl : baseUrl + "/" + keyspace; - // Multi-DC clusters require localdatacenter= String extraParams = coalesce(params.getUrl_params()).trim(); while (extraParams.startsWith("?") || extraParams.startsWith("&")) { extraParams = extraParams.substring(1); } + if (!hasDatacenterOrLoadBalancingParameter(extraParams)) { + extraParams = extraParams.isEmpty() + ? DC_INFERRING_LOAD_BALANCING + : extraParams + (extraParams.endsWith("&") ? "" : "&") + DC_INFERRING_LOAD_BALANCING; + } return extraParams.isEmpty() ? url : url + "?" + extraParams; } + private static boolean hasDatacenterOrLoadBalancingParameter(String urlParams) { + for (String param : urlParams.split("&")) { + int valueSeparator = param.indexOf('='); + String paramName = valueSeparator < 0 ? param : param.substring(0, valueSeparator); + paramName = paramName.trim(); + if (paramName.equals("localdatacenter") || paramName.equals("loadbalancing")) { + return true; + } + } + return false; + } + private static List targetColumns(String options) { Matcher matcher = TARGET_PATTERN.matcher(options); if (!matcher.find()) { diff --git a/agents/drivers/cassandra/src/test/java/com/dbx/agent/cassandra/CassandraAgentTest.java b/agents/drivers/cassandra/src/test/java/com/dbx/agent/cassandra/CassandraAgentTest.java index 4eb9e6be6..5873d8abc 100644 --- a/agents/drivers/cassandra/src/test/java/com/dbx/agent/cassandra/CassandraAgentTest.java +++ b/agents/drivers/cassandra/src/test/java/com/dbx/agent/cassandra/CassandraAgentTest.java @@ -22,14 +22,20 @@ class CassandraAgentTest extends JdbcFakeExecutionBehaviorTest { void buildsServerUrlWhenKeyspaceIsEmpty() { ConnectParams params = new ConnectParams("127.0.0.1", 9042, "", "cassandra", "cassandra", "", "", false); - assertEquals("jdbc:cassandra://127.0.0.1:9042", CassandraAgent.buildUrl(params)); + assertEquals( + "jdbc:cassandra://127.0.0.1:9042?loadbalancing=DcInferringLoadBalancingPolicy", + CassandraAgent.buildUrl(params) + ); } @Test void buildsKeyspaceUrlWhenKeyspaceIsSet() { ConnectParams params = new ConnectParams("127.0.0.1", 9042, "app_keyspace", "cassandra", "cassandra", "", "", false); - assertEquals("jdbc:cassandra://127.0.0.1:9042/app_keyspace", CassandraAgent.buildUrl(params)); + assertEquals( + "jdbc:cassandra://127.0.0.1:9042/app_keyspace?loadbalancing=DcInferringLoadBalancingPolicy", + CassandraAgent.buildUrl(params) + ); } @Test @@ -49,4 +55,40 @@ class CassandraAgentTest extends JdbcFakeExecutionBehaviorTest { assertEquals("jdbc:cassandra://127.0.0.1:9042?localdatacenter=dc1", CassandraAgent.buildUrl(params)); } + + @Test + void appendsInferringPolicyAfterOtherUrlParams() { + ConnectParams params = new ConnectParams( + "127.0.0.1", 9042, "", "cassandra", "cassandra", "requesttimeout=10000", "", false + ); + + assertEquals( + "jdbc:cassandra://127.0.0.1:9042?requesttimeout=10000&loadbalancing=DcInferringLoadBalancingPolicy", + CassandraAgent.buildUrl(params) + ); + } + + @Test + void appendsInferringPolicyAfterTrailingSeparator() { + ConnectParams params = new ConnectParams( + "127.0.0.1", 9042, "", "cassandra", "cassandra", "requesttimeout=10000&", "", false + ); + + assertEquals( + "jdbc:cassandra://127.0.0.1:9042?requesttimeout=10000&loadbalancing=DcInferringLoadBalancingPolicy", + CassandraAgent.buildUrl(params) + ); + } + + @Test + void preservesCustomLoadBalancingPolicy() { + ConnectParams params = new ConnectParams( + "127.0.0.1", 9042, "", "cassandra", "cassandra", "loadbalancing=CustomPolicy", "", false + ); + + assertEquals( + "jdbc:cassandra://127.0.0.1:9042?loadbalancing=CustomPolicy", + CassandraAgent.buildUrl(params) + ); + } }