fix(cassandra): infer local datacenter when unspecified

Closes #976
This commit is contained in:
t8y2 2026-08-03 18:04:00 +08:00
parent 892a40274e
commit e4d22904ed
No known key found for this signature in database
2 changed files with 62 additions and 3 deletions

View File

@ -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=<dc>
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<String> targetColumns(String options) {
Matcher matcher = TARGET_PATTERN.matcher(options);
if (!matcher.find()) {

View File

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