From b828f7ea1b735586da388ddfee17f26685e20cef Mon Sep 17 00:00:00 2001 From: Stefan Miklosovic Date: Thu, 18 May 2023 11:32:18 +0200 Subject: [PATCH] Pass down all contact points to driver for cassandra-stress patch by Stefan Miklosovic; reviewed by Brandon Williams for CASSANDRA-18025 --- CHANGES.txt | 1 + .../stress/settings/StressSettings.java | 3 +-- .../stress/util/JavaDriverClient.java | 24 ++++++++++++++----- 3 files changed, 20 insertions(+), 8 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index a9fd724d6a..fe4b307e1b 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0.30 + * Pass down all contact points to driver for cassandra-stress (CASSANDRA-18025) * Validate the existence of a datacenter in nodetool rebuild (CASSANDRA-14319) 3.0.29 diff --git a/tools/stress/src/org/apache/cassandra/stress/settings/StressSettings.java b/tools/stress/src/org/apache/cassandra/stress/settings/StressSettings.java index 5b1f8611cf..f16f1aa930 100644 --- a/tools/stress/src/org/apache/cassandra/stress/settings/StressSettings.java +++ b/tools/stress/src/org/apache/cassandra/stress/settings/StressSettings.java @@ -182,12 +182,11 @@ public class StressSettings implements Serializable { synchronized (this) { - String currentNode = node.randomNode(); if (client != null) return client; EncryptionOptions.ClientEncryptionOptions encOptions = transport.getEncryptionOptions(); - JavaDriverClient c = new JavaDriverClient(this, currentNode, port.nativePort, encOptions); + JavaDriverClient c = new JavaDriverClient(this, node.nodes, port.nativePort, encOptions); c.connect(mode.compression()); if (setKeyspace) c.execute("USE \"" + schema.keyspace + "\";", org.apache.cassandra.db.ConsistencyLevel.ONE); diff --git a/tools/stress/src/org/apache/cassandra/stress/util/JavaDriverClient.java b/tools/stress/src/org/apache/cassandra/stress/util/JavaDriverClient.java index 4f173b4628..2720d969cc 100644 --- a/tools/stress/src/org/apache/cassandra/stress/util/JavaDriverClient.java +++ b/tools/stress/src/org/apache/cassandra/stress/util/JavaDriverClient.java @@ -17,11 +17,16 @@ */ package org.apache.cassandra.stress.util; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import javax.net.ssl.SSLContext; +import com.google.common.net.HostAndPort; import com.datastax.driver.core.*; import com.datastax.driver.core.policies.DCAwareRoundRobinPolicy; import com.datastax.driver.core.policies.WhiteListPolicy; @@ -39,7 +44,7 @@ public class JavaDriverClient InternalLoggerFactory.setDefaultFactory(new Slf4JLoggerFactory()); } - public final String host; + public final List hosts; public final int port; public final String username; public final String password; @@ -57,13 +62,13 @@ public class JavaDriverClient public JavaDriverClient(StressSettings settings, String host, int port) { - this(settings, host, port, new EncryptionOptions.ClientEncryptionOptions()); + this(settings, Collections.singletonList(host), port, new EncryptionOptions.ClientEncryptionOptions()); } - public JavaDriverClient(StressSettings settings, String host, int port, EncryptionOptions.ClientEncryptionOptions encryptionOptions) + public JavaDriverClient(StressSettings settings, List hosts, int port, EncryptionOptions.ClientEncryptionOptions encryptionOptions) { this.protocolVersion = settings.mode.protocolVersion; - this.host = host; + this.hosts = hosts; this.port = port; this.username = settings.mode.username; this.password = settings.mode.password; @@ -112,9 +117,16 @@ public class JavaDriverClient .setMaxRequestsPerConnection(HostDistance.LOCAL, maxPendingPerConnection) .setNewConnectionThreshold(HostDistance.LOCAL, 100); + List contacts = new ArrayList<>(); + for (String host : hosts) + { + HostAndPort hap = HostAndPort.fromString(host).withDefaultPort(port); + InetSocketAddress contact = new InetSocketAddress(InetAddress.getByName(hap.getHostText()), hap.getPort()); + contacts.add(contact); + } + Cluster.Builder clusterBuilder = Cluster.builder() - .addContactPoint(host) - .withPort(port) + .addContactPointsWithPorts(contacts) .withPoolingOptions(poolingOpts) .withoutJMXReporting() .withProtocolVersion(protocolVersion)