From c52fb9b70d1d5a165eb7c0b462bfdf2ce95aaa55 Mon Sep 17 00:00:00 2001 From: Eric Pugh Date: Sat, 15 Aug 2026 18:11:07 -0400 Subject: [PATCH 1/3] Tidy up some code for issues marked by IntelliJ --- .../manager/consumer/ConsumerMetrics.java | 2 +- .../manager/consumer/PartitionManager.java | 16 -------- .../crossdc/manager/consumer/ThreadDump.java | 6 +-- .../manager/consumer/ThreadDumpServlet.java | 2 +- .../SolrMessageProcessor.java | 41 +------------------ .../manager/DeleteByQueryToIdTest.java | 3 -- .../manager/RetryQueueIntegrationTest.java | 12 ++---- .../manager/SolrAndKafkaIntegrationTest.java | 6 +-- ...ndKafkaMultiCollectionIntegrationTest.java | 5 +-- .../manager/SolrAndKafkaReindexTest.java | 7 +--- .../manager/ZkConfigIntegrationTest.java | 8 +--- .../consumer/KafkaCrossDcConsumerTest.java | 27 +----------- .../TestMessageProcessor.java | 2 +- 13 files changed, 21 insertions(+), 116 deletions(-) diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java index 92d96fae866e..7d5f11a427eb 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ConsumerMetrics.java @@ -131,7 +131,7 @@ default void incrementOutputCounter(String type, String result) { /** * Records the batch size of the output request. Batch size is defined as the number of operations - * in an output {@link SolrRequest} (which may be different than the input size due to + * in an output {@link SolrRequest} (which may be different from the input size due to * collapsing). * * @param type the type of the request, corresponding to one of the {@link diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java index c93740f25ea0..cef9bfefb986 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/PartitionManager.java @@ -119,22 +119,6 @@ void checkForOffsetUpdates(TopicPartition partition) throws Throwable { } } - /** - * Reset the local offset so that the consumer reads the records from Kafka again. - * - * @param partition The TopicPartition to reset the offset for - * @param partitionRecords PartitionRecords for the specified partition - */ - private void resetOffsetForPartition( - TopicPartition partition, - List>> partitionRecords) { - if (log.isTraceEnabled()) { - log.trace("Resetting offset to: {}", partitionRecords.get(0).offset()); - } - long resetOffset = partitionRecords.get(0).offset(); - consumer.seek(partition, resetOffset); - } - /** * Logs and updates the commit point for the partition that has been processed. * diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java index 2f0fb116bf9d..aed28f2c9528 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDump.java @@ -42,7 +42,7 @@ public ThreadDump(ThreadMXBean threadMXBean) { } /** - * Dumps all of the threads' current information, including synchronization, to an output stream. + * Dumps all the threads' current information, including synchronization, to an output stream. * * @param out an output stream */ @@ -51,8 +51,8 @@ public void dump(OutputStream out) { } /** - * Dumps all of the threads' current information, optionally including synchronization, to an - * output stream. + * Dumps all the threads' current information, optionally including synchronization, to an output + * stream. * *

Having control over including synchronization info allows using this method (and its * wrappers, i.e. ThreadDumpServlet) in environments where getting object monitor and/or ownable diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java index cea2bb339884..7df3e16b0e22 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/consumer/ThreadDumpServlet.java @@ -45,7 +45,7 @@ public class ThreadDumpServlet extends HttpServlet { @Override public void init() throws ServletException { try { - // Some PaaS like Google App Engine blacklist java.lang.managament + // Some PaaS like Google App Engine blacklist java.lang.management this.threadDump = new ThreadDump(ManagementFactory.getThreadMXBean()); } catch (NoClassDefFoundError ncdfe) { this.threadDump = null; // we won't be able to provide thread dump diff --git a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java index de3484ca45a6..af9f7501ee6b 100644 --- a/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java +++ b/solr/cross-dc-manager/src/java/org/apache/solr/crossdc/manager/messageprocessor/SolrMessageProcessor.java @@ -31,10 +31,8 @@ import org.apache.solr.common.cloud.ClusterState; import org.apache.solr.common.cloud.DocCollection; import org.apache.solr.common.cloud.ZkStateReader; -import org.apache.solr.common.params.ModifiableSolrParams; import org.apache.solr.common.params.SolrParams; import org.apache.solr.common.util.TimeSource; -import org.apache.solr.crossdc.common.CrossDcConstants; import org.apache.solr.crossdc.common.IQueueHandler; import org.apache.solr.crossdc.common.MirroredSolrRequest; import org.apache.solr.crossdc.common.ResubmitBackoffPolicy; @@ -165,7 +163,7 @@ private boolean isRetryable(Exception e) { } private void logIf4xxException(SolrException solrException) { - // This shouldn't really happen but if it doesn, it most likely requires fixing in the return + // This shouldn't really happen but if it does, it most likely requires fixing in the return // code from Solr. if (solrException != null && 400 <= solrException.code() && solrException.code() < 500) { log.error("Exception occurred with 4xx response. {}", solrException.code(), solrException); @@ -335,43 +333,6 @@ private void logFirstAttemptLatency(MirroredSolrRequest mirroredSolrRequest) } } - /** - * Adds {@link CrossDcConstants#SHOULD_MIRROR}=false to the params if it's not already specified. - * Logs a warning if it is specified and NOT set to false. (i.e. circular mirror may occur) - * - * @param mirroredSolrRequest MirroredSolrRequest object that is being processed. - */ - void preventCircularMirroring(MirroredSolrRequest mirroredSolrRequest) { - if (mirroredSolrRequest.getSolrRequest() instanceof UpdateRequest updateRequest) { - ModifiableSolrParams params = updateRequest.getParams(); - String shouldMirror = (params == null ? null : params.get(CrossDcConstants.SHOULD_MIRROR)); - if (shouldMirror == null) { - log.warn( - "{} param is missing - setting to false. Request={}", - CrossDcConstants.SHOULD_MIRROR, - mirroredSolrRequest); - updateRequest.setParam(CrossDcConstants.SHOULD_MIRROR, "false"); - } else if (!"false".equalsIgnoreCase(shouldMirror)) { - log.warn("{} param equal to {}", CrossDcConstants.SHOULD_MIRROR, shouldMirror); - } - } else { - SolrParams params = mirroredSolrRequest.getSolrRequest().getParams(); - assert params != null; - String shouldMirror = params.get(CrossDcConstants.SHOULD_MIRROR); - if (shouldMirror == null) { - if (params instanceof ModifiableSolrParams) { - log.warn("{} param is missing - setting to false", CrossDcConstants.SHOULD_MIRROR); - ((ModifiableSolrParams) params).set(CrossDcConstants.SHOULD_MIRROR, "false"); - } else { - log.warn( - "{} param is missing and params are not modifiable", CrossDcConstants.SHOULD_MIRROR); - } - } else if (!"false".equalsIgnoreCase(shouldMirror)) { - log.warn("{} param is present and set to {}", CrossDcConstants.SHOULD_MIRROR, shouldMirror); - } - } - } - private void connectToSolrIfNeeded() { // Don't try to consume anything if we can't connect to the solr server boolean connected = false; diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java index 1898def635cd..b73b6333984b 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/DeleteByQueryToIdTest.java @@ -55,7 +55,6 @@ import org.slf4j.LoggerFactory; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, @@ -67,8 +66,6 @@ public class DeleteByQueryToIdTest extends SolrCloudTestCase { private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - static final String VERSION_FIELD = "_version_"; - private static final int NUM_BROKERS = 1; public static EmbeddedKafkaCluster kafkaCluster; diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java index 1e06264c965a..dca1b60d09ba 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/RetryQueueIntegrationTest.java @@ -60,8 +60,6 @@ public class RetryQueueIntegrationTest extends SolrTestCaseJ4 { private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - static final String VERSION_FIELD = "_version_"; - private static final int NUM_BROKERS = 1; public static EmbeddedKafkaCluster kafkaCluster; @@ -127,10 +125,8 @@ public String bootstrapServers() { zkTestServer2.run(); } - solrCluster1 = startCluster(solrCluster1, zkTestServer1, baseDir1); - solrCluster2 = startCluster(solrCluster2, zkTestServer2, baseDir2); - - CloudSolrClient client = solrCluster1.getSolrClient(COLLECTION); + solrCluster1 = startCluster(zkTestServer1, baseDir1); + solrCluster2 = startCluster(zkTestServer2, baseDir2); String bootstrapServers = kafkaCluster.bootstrapServers(); log.info("bootstrapServers={}", bootstrapServers); @@ -144,8 +140,8 @@ public String bootstrapServers() { consumer.start(properties); } - private static MiniSolrCloudCluster startCluster( - MiniSolrCloudCluster solrCluster, ZkTestServer zkTestServer, Path baseDir) throws Exception { + private static MiniSolrCloudCluster startCluster(ZkTestServer zkTestServer, Path baseDir) + throws Exception { MiniSolrCloudCluster cluster = new MiniSolrCloudCluster( 1, diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java index 797080454b5a..9a19671feb24 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaIntegrationTest.java @@ -586,14 +586,14 @@ public void testMetricsAndHealthcheck() throws Exception { (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), StandardCharsets.UTF_8); assertTrue(content, content.contains("solr_crossdc_consumer_output_total")); - // test the healtcheck endpoint + // test the healthcheck endpoint req = new GenericSolrRequest(SolrRequest.METHOD.GET, "/health"); req.setResponseParser(new InputStreamResponseParser(null)); rsp = httpJettySolrClient.request(req); content = IOUtils.toString( (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), StandardCharsets.UTF_8); - assertEquals(Integer.valueOf(200), rsp.get("responseStatus")); + assertEquals(200, rsp.get("responseStatus")); Map map = (Map) ObjectBuilder.fromJSON(content); assertEquals(Boolean.TRUE, map.get("kafka")); assertEquals(Boolean.TRUE, map.get("solr")); @@ -607,7 +607,7 @@ public void testMetricsAndHealthcheck() throws Exception { content = IOUtils.toString( (InputStream) rsp.get(InputStreamResponseParser.STREAM_KEY), StandardCharsets.UTF_8); - assertEquals(Integer.valueOf(503), rsp.get("responseStatus")); + assertEquals(503, rsp.get("responseStatus")); map = (Map) ObjectBuilder.fromJSON(content); assertEquals(Boolean.TRUE, map.get("kafka")); assertEquals(Boolean.FALSE, map.get("solr")); diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java index e859e80f4b88..917a9e2af834 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaMultiCollectionIntegrationTest.java @@ -58,22 +58,19 @@ import org.slf4j.LoggerFactory; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, SolrKafkaTestsIgnoredThreadsFilter.class }) @ThreadLeakLingering(linger = 5000) -@Ignore("This test relies on collecton properties and I don't see where they are set anymore") +@Ignore("This test relies on collection properties and I don't see where they are set anymore") public class SolrAndKafkaMultiCollectionIntegrationTest extends SolrCloudTestCase { private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); private static final int MAX_DOC_SIZE_BYTES = Integer.parseInt(DEFAULT_MAX_REQUEST_SIZE); - static final String VERSION_FIELD = "_version_"; - private static final int NUM_BROKERS = 1; public EmbeddedKafkaCluster kafkaCluster; diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java index a7c12207e57f..b01163d8061c 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/SolrAndKafkaReindexTest.java @@ -47,7 +47,6 @@ import org.slf4j.LoggerFactory; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, @@ -58,8 +57,6 @@ public class SolrAndKafkaReindexTest extends SolrCloudTestCase { private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - static final String VERSION_FIELD = "_version_"; - private static final int NUM_BROKERS = 1; public static EmbeddedKafkaCluster kafkaCluster; @@ -126,7 +123,7 @@ public String bootstrapServers() { } @AfterClass - public static void afterSolrAndKafkaIntegrationTest() throws Exception { + public static void afterSolrAndKafkaIntegrationTest() { ObjectReleaseTracker.clear(); if (solrCluster1 != null) { @@ -267,7 +264,7 @@ private void addDocs(CloudSolrClient client, String tag) throws SolrServerExcept doc2.addField("id", id2); doc2.addField("text", "some test two " + tag); - List docs = new ArrayList(2); + List docs = new ArrayList<>(2); docs.add(doc1); docs.add(doc2); diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java index e0af7eaeb612..aa55a88cd3fa 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/ZkConfigIntegrationTest.java @@ -45,7 +45,6 @@ import org.slf4j.LoggerFactory; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, @@ -57,8 +56,6 @@ public class ZkConfigIntegrationTest extends SolrCloudTestCase { private static final Logger log = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - static final String VERSION_FIELD = "_version_"; - private static final int NUM_BROKERS = 1; public static EmbeddedKafkaCluster kafkaCluster; @@ -140,9 +137,8 @@ public String bootstrapServers() { log.info("bootstrapServers={}", bootstrapServers); Map properties = new HashMap<>(); - Object put = - properties.put( - KafkaCrossDcConf.ZK_CONNECT_STRING, solrCluster2.getZkServer().getZkAddress()); + + properties.put(KafkaCrossDcConf.ZK_CONNECT_STRING, solrCluster2.getZkServer().getZkAddress()); System.setProperty(KafkaCrossDcConf.BOOTSTRAP_SERVERS, kafkaCluster.bootstrapServers()); diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java index fcbd2edd84d1..2bc50db48fd9 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/consumer/KafkaCrossDcConsumerTest.java @@ -151,29 +151,6 @@ public void tearDown() { kafkaCrossDcConsumer.shutdown(); } - private ConsumerRecord> createSampleConsumerRecord() { - return new ConsumerRecord<>("sample-topic", 0, 0, "key", createSampleMirroredSolrRequest()); - } - - private ConsumerRecords> createSampleConsumerRecords() { - TopicPartition topicPartition = new TopicPartition("sample-topic", 0); - List>> recordsList = new ArrayList<>(); - recordsList.add( - new ConsumerRecord<>("sample-topic", 0, 0, "key", createSampleMirroredSolrRequest())); - return new ConsumerRecords<>(Map.of(topicPartition, recordsList)); - } - - private MirroredSolrRequest createSampleMirroredSolrRequest() { - // Create a sample MirroredSolrRequest for testing - SolrInputDocument solrInputDocument = new SolrInputDocument(); - solrInputDocument.addField("id", "1"); - solrInputDocument.addField("title", "Sample title"); - solrInputDocument.addField("content", "Sample content"); - UpdateRequest updateRequest = new UpdateRequest(); - updateRequest.add(solrInputDocument); - return new MirroredSolrRequest<>(updateRequest); - } - /** Should create a KafkaCrossDcConsumer with the given configuration and startLatch */ @Test public void kafkaCrossDcConsumerCreationWithConfigurationAndStartLatch() { @@ -204,7 +181,7 @@ protected KafkaMirroringSink createKafkaMirroringSink(KafkaCrossDcConf conf) { } @Test - public void testSolrClientSupplier() throws Exception { + public void testSolrClientSupplier() { supplier.get(); assertEquals(1, solrClientCounter.get()); clusterStateProviderIsClosed = true; @@ -331,7 +308,7 @@ public void testHandleValidMirroredSolrRequest() { } @Test - public void testHandleValidAdminRequest() throws Exception { + public void testHandleValidAdminRequest() { KafkaConsumer> mockConsumer = mock(KafkaConsumer.class); KafkaCrossDcConsumer spyConsumer = createCrossDcConsumerSpy(mockConsumer); doReturn(new IQueueHandler.Result<>(IQueueHandler.ResultStatus.HANDLED, null)) diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java index 769fc58cb29c..5cb94f94dec8 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java @@ -70,7 +70,7 @@ public static void ensureWorkingMockito() { @Before public void setUp() { - MockitoAnnotations.initMocks(this); + MockitoAnnotations.openMocks(this); ConsumerMetrics metrics = Mockito.mock(OtelMetrics.class); processor = Mockito.spy(new SolrMessageProcessor(metrics, () -> solrClient, backoffPolicy)); From f406132ed21f6c2de1ae7a8f01395d9f693266c1 Mon Sep 17 00:00:00 2001 From: Eric Pugh Date: Mon, 17 Aug 2026 11:45:50 -0400 Subject: [PATCH 2/3] tidy up the crossdc module code for Solr --- .../solr/crossdc/common/ConfigProperty.java | 10 ------ .../solr/crossdc/common/IQueueHandler.java | 2 +- .../solr/crossdc/common/KafkaCrossDcConf.java | 6 ++-- .../crossdc/common/KafkaMirroringSink.java | 5 ++- .../crossdc/common/MirroredSolrRequest.java | 21 +++++------ .../update/processor/MirroringException.java | 36 ------------------- .../processor/MirroringUpdateProcessor.java | 10 ++---- ...irroringUpdateRequestProcessorFactory.java | 15 +------- .../solr/crossdc/common/ConfUtilTest.java | 2 +- .../MirroredSolrRequestSerializerTest.java | 4 +-- .../MirroringCollectionsHandlerTest.java | 3 +- .../MirroringConfigSetsHandlerTest.java | 6 ++-- .../MirroringUpdateProcessorTest.java | 14 ++++---- 13 files changed, 31 insertions(+), 103 deletions(-) delete mode 100644 solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java index fc2bd416ec66..660f8ce2d944 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/ConfigProperty.java @@ -25,12 +25,6 @@ public class ConfigProperty { private boolean required = false; - public ConfigProperty(String key, String defaultValue, boolean required) { - this.key = key; - this.defaultValue = defaultValue; - this.required = required; - } - public ConfigProperty(String key, String defaultValue) { this.key = key; this.defaultValue = defaultValue; @@ -49,10 +43,6 @@ public boolean isRequired() { return required; } - public String getDefaultValue() { - return defaultValue; - } - public String getValue(Map properties) { String val = (String) properties.get(key); if (val == null) { diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java index 145630deee7a..c7f852e3716f 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/IQueueHandler.java @@ -21,7 +21,7 @@ enum ResultStatus { /** Item was successfully processed */ HANDLED, - /** Item was not processed, and the consumer should shutdown */ + /** Item was not processed, and the consumer should shut down */ NOT_HANDLED_SHUTDOWN, /** Item processing failed, and the item should be retried immediately */ diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java index 02526216b36e..88db372683e9 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaCrossDcConf.java @@ -225,7 +225,7 @@ public class KafkaCrossDcConf extends CrossDcConf { private final Map properties; public KafkaCrossDcConf(Map properties) { - List nullValueKeys = new ArrayList(); + List nullValueKeys = new ArrayList<>(); properties.forEach( (k, v) -> { if (v == null) { @@ -275,7 +275,7 @@ public Map getAdditionalProperties() { (key, v) -> { try { int intVal = Integer.parseInt((String) v); - integerProperties.put(key.toString(), intVal); + integerProperties.put(key, intVal); } catch (NumberFormatException ignored) { } @@ -312,7 +312,7 @@ public String toString() { sb.append(configProperty.getKey()).append("=").append(printablePropertyValue).append(","); } } - if (sb.length() > 0) { + if (!sb.isEmpty()) { sb.setLength(sb.length() - 1); } diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java index 70f484578145..c348f032c21d 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/KafkaMirroringSink.java @@ -110,13 +110,12 @@ private void submitRequest(MirroredSolrRequest request, String topicName) } }); - long lastSuccessfulEnqueueNanos = System.nanoTime(); // Record time since last successful enqueue as 0 long elapsedTimeMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - enqueueStartNanos); // Update elapsed time if (elapsedTimeMillis > conf.getInt(SLOW_SUBMIT_THRESHOLD_MS)) { - slowSubmitAction(request, elapsedTimeMillis); + slowSubmitAction(elapsedTimeMillis); } } catch (Exception e) { // We are intentionally catching all exceptions, the expected exception form this function is @@ -220,7 +219,7 @@ private KafkaConsumer> initConsumer() { kafkaConsumerProperties, new StringDeserializer(), new MirroredSolrRequestSerializer()); } - private void slowSubmitAction(Object request, long elapsedTimeMillis) { + private void slowSubmitAction(long elapsedTimeMillis) { log.warn( "Enqueuing the request to Kafka took more than {} millis. enqueueElapsedTime={}", conf.get(KafkaCrossDcConf.SLOW_SUBMIT_THRESHOLD_MS), diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java index 7d27878f89b9..4d29239adbe8 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/common/MirroredSolrRequest.java @@ -49,7 +49,7 @@ public enum Type { CONFIGSET, UNKNOWN; - public static final Type get(String s) { + public static Type get(String s) { if (s == null) { return UNKNOWN; } else { @@ -227,10 +227,6 @@ public long getSubmitTimeNanos() { return submitTimeNanos; } - public void setSubmitTimeNanos(final long submitTimeNanos) { - this.submitTimeNanos = submitTimeNanos; - } - public Type getType() { return type; } @@ -264,15 +260,16 @@ public int hashCode() { public String toString() { final StringBuilder sb = new StringBuilder(getClass().getSimpleName() + "{type="); sb.append(type.toString()); - sb.append(", method=" + solrRequest.getMethod()); - sb.append(", params=" + solrRequest.getParams()); + sb.append(", method=").append(solrRequest.getMethod()); + sb.append(", params=").append(solrRequest.getParams()); if (solrRequest instanceof UpdateRequest req) { - sb.append(", add=" + (req.getDocuments() != null ? req.getDocuments().size() : "0")); - sb.append(", del=" + (req.getDeleteByIdMap() != null ? req.getDeleteByIdMap().size() : "0")); - sb.append(", dbq=" + (req.getDeleteQuery() != null ? req.getDeleteQuery().size() : "0")); + sb.append(", add=").append(req.getDocuments() != null ? req.getDocuments().size() : "0"); + sb.append(", del=") + .append(req.getDeleteByIdMap() != null ? req.getDeleteByIdMap().size() : "0"); + sb.append(", dbq=").append(req.getDeleteQuery() != null ? req.getDeleteQuery().size() : "0"); } - sb.append(", attempt=" + attempt); - sb.append(", submitTimeNanos=" + submitTimeNanos); + sb.append(", attempt=").append(attempt); + sb.append(", submitTimeNanos=").append(submitTimeNanos); sb.append('}'); return sb.toString(); } diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java deleted file mode 100644 index cf01fa838248..000000000000 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringException.java +++ /dev/null @@ -1,36 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.solr.crossdc.update.processor; - -/** Wrapper class for Mirroring exceptions. */ -public class MirroringException extends Exception { - public MirroringException() { - super(); - } - - public MirroringException(String message) { - super(message); - } - - public MirroringException(String message, Throwable cause) { - super(message, cause); - } - - public MirroringException(Throwable cause) { - super(cause); - } -} diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java index 54ad07e3d243..284dd91ca7d7 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessor.java @@ -86,17 +86,11 @@ public class MirroringUpdateProcessor extends UpdateRequestProcessor { /** If true then commit commands are mirrored, otherwise they are processed only locally. */ private final boolean mirrorCommits; - /** Controls the processing of Delete-By-Query requests.. */ + /** Controls the processing of Delete-By-Query requests */ private final CrossDcConf.ExpandDbq expandDbq; private final long maxMirroringDocSizeBytes; - /** - * The distributed processor downstream from us so we can establish if we're running on a leader - * shard - */ - // private DistributedUpdateProcessor distProc; - /** Distribution phase of the incoming requests */ private DistributedUpdateProcessor.DistribPhase distribPhase; @@ -255,7 +249,7 @@ public void processDelete(final DeleteUpdateCommand cmd) throws IOException { super.processDelete(cmd); // let this throw to prevent mirroring invalid requests producerMetrics.getLocal().inc(); if (doMirroring) { - boolean isLeader = false; + boolean isLeader; UpdateRequest mirrorRequest = createMirrorRequest(); if (cmd.isDeleteById()) { // deleteById requests runs once per leader, so we just submit the request from the leader diff --git a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java index c4329b776f78..b159d060fc00 100644 --- a/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java +++ b/solr/modules/cross-dc/src/java/org/apache/solr/crossdc/update/processor/MirroringUpdateRequestProcessorFactory.java @@ -34,7 +34,6 @@ import java.lang.invoke.MethodHandles; import java.util.HashMap; import java.util.Map; -import java.util.Properties; import org.apache.solr.common.SolrException; import org.apache.solr.common.cloud.CollectionProperties; import org.apache.solr.common.cloud.SolrZkClient; @@ -165,11 +164,7 @@ private void lookupPropertyOverridesInZk(SolrCore core) { } String enabledVal = collectionProperties.get("crossdc.enabled"); if (enabledVal != null) { - if (Boolean.parseBoolean(enabledVal.toString())) { - this.enabled = true; - } else { - this.enabled = false; - } + this.enabled = Boolean.parseBoolean(enabledVal); } } catch (Exception e) { log.error("Exception looking for CrossDC configuration in Zookeeper", e); @@ -223,14 +218,6 @@ public void inform(SolrCore core) { mirroringHandler = new KafkaRequestMirroringHandler(sink); } - private static Integer getIntegerPropValue(String name, Properties props) { - String value = props.getProperty(name); - if (value == null) { - return null; - } - return Integer.parseInt(value); - } - @Override public UpdateRequestProcessor getInstance( final SolrQueryRequest req, final SolrQueryResponse rsp, final UpdateRequestProcessor next) { diff --git a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java index 28a98d6d7a1e..53713cdfc122 100644 --- a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java +++ b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/ConfUtilTest.java @@ -291,7 +291,7 @@ public void testFillProperties_EmptyProperties() { } @Test - public void testFillProperties_SecurityProperties() throws Exception { + public void testFillProperties_SecurityProperties() { Map properties = new HashMap<>(); // Set security-related properties diff --git a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java index 2ce667fee641..e0afa71cd33f 100644 --- a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java +++ b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/common/MirroredSolrRequestSerializerTest.java @@ -28,7 +28,7 @@ public class MirroredSolrRequestSerializerTest extends SolrTestCase { private static final byte[] EMPTY_ARR = new byte[3]; @Test - public void testSerializationBufferOptimization() throws Exception { + public void testSerializationBufferOptimization() { MirroredSolrRequestSerializer serializer = new MirroredSolrRequestSerializer(); UpdateRequest req = new UpdateRequest(); SolrInputDocument doc = new SolrInputDocument(); @@ -50,7 +50,7 @@ public void testSerializationBufferOptimization() throws Exception { (String) ((UpdateRequest) deserialized.getSolrRequest()) .getDocuments() - .get(0) + .getFirst() .getFieldValue("test"); assertEquals(fieldValue, deserValue); } diff --git a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java index d2a7dc33df23..a2d0537d57d2 100644 --- a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java +++ b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringCollectionsHandlerTest.java @@ -46,7 +46,6 @@ import org.mockito.Mockito; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, @@ -141,7 +140,7 @@ private void runCommand(SolrParams params, boolean expectResult) throws Exceptio SolrParams mirroredParams = solrRequest.getParams(); params.forEach( entry -> { - assertEquals(entry.getValue(), mirroredParams.getParams(entry.getKey())); + assertArrayEquals(entry.getValue(), mirroredParams.getParams(entry.getKey())); }); } else { assertEquals(initialMirroredCount, captor.getAllValues().size()); diff --git a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java index 275fa528385a..b115f4cdcdbc 100644 --- a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java +++ b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/handler/MirroringConfigSetsHandlerTest.java @@ -20,7 +20,6 @@ import com.carrotsearch.randomizedtesting.annotations.ThreadLeakLingering; import java.nio.charset.StandardCharsets; import java.nio.file.Path; -import java.util.Arrays; import java.util.List; import org.apache.commons.io.IOUtils; import org.apache.lucene.tests.util.QuickPatchThreadsFilter; @@ -52,7 +51,6 @@ import org.mockito.Mockito; @ThreadLeakFilters( - defaultFilters = true, filters = { SolrIgnoredThreadsFilter.class, QuickPatchThreadsFilter.class, @@ -149,7 +147,7 @@ private void runCommand(SolrQueryRequest req, boolean expectStreams, boolean exp req.getParams() .forEach( entry -> { - assertEquals(entry.getValue(), mirroredParams.getParams(entry.getKey())); + assertArrayEquals(entry.getValue(), mirroredParams.getParams(entry.getKey())); }); assertEquals("HTTP method", req.getHttpMethod(), solrRequest.getMethod().toString()); if (expectStreams) { @@ -167,7 +165,7 @@ private void runCommand(SolrQueryRequest req, boolean expectStreams, boolean exp MirroredSolrRequest.ExposedByteArrayContentStream.of(source).byteArray(); byte[] mirroredContent = MirroredSolrRequest.ExposedByteArrayContentStream.of(mirrored).byteArray(); - assertTrue("different content", Arrays.equals(sourceContent, mirroredContent)); + assertArrayEquals("different content", sourceContent, mirroredContent); } } } else { diff --git a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java index 77b3adf7a494..476b4529001c 100644 --- a/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java +++ b/solr/modules/cross-dc/src/test/org/apache/solr/crossdc/update/processor/MirroringUpdateProcessorTest.java @@ -544,7 +544,7 @@ UpdateRequest createMirrorRequest() { processor.processDelete(deleteUpdateCommand); verify(requestMirroringHandler, times(1)).mirror(updateRequest); assertEquals("missing dbq", 1, updateRequest.getDeleteQuery().size()); - assertEquals("dbq value", "id:test*", updateRequest.getDeleteQuery().get(0)); + assertEquals("dbq value", "id:test*", updateRequest.getDeleteQuery().getFirst()); // verify the metrics assertEquals(1, counters.get("local").get()); assertEquals(1, counters.get("submittedDeleteByQuery").get()); @@ -575,13 +575,13 @@ public void testProcessDBQResults() throws Exception { @Test public void testEstimateObjectSize() { - assertEquals(estimate(null), 0); - assertEquals(estimate("abc"), 6); - assertEquals(estimate("abcdefgh"), 16); + assertEquals(0, estimate(null)); + assertEquals(6, estimate("abc")); + assertEquals(16, estimate("abcdefgh")); List keys = List.of("int", "long", "double", "float", "str"); - assertEquals(estimate(keys), 42); + assertEquals(42, estimate(keys)); List values = List.of(12, 5L, 12.0, 5.0, "duck"); - assertEquals(estimate(values), 8); + assertEquals(8, estimate(values)); Map map = new HashMap<>(); map.put("int", 12); @@ -590,7 +590,7 @@ public void testEstimateObjectSize() { map.put("float", 5.0f); map.put("str", "duck"); map.put("short", null); - assertEquals(estimate(map), 60); + assertEquals(60, estimate(map)); SolrInputDocument document = new SolrInputDocument(); for (Map.Entry entry : map.entrySet()) { From b311a3535a21c7221b5c4ebb0e63e7876bd90f43 Mon Sep 17 00:00:00 2001 From: Eric Pugh Date: Tue, 18 Aug 2026 16:19:35 -0400 Subject: [PATCH 3/3] Follow pattern in ConfUtilTest. Could be a rule someday? --- .../messageprocessor/TestMessageProcessor.java | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java index 5cb94f94dec8..c6c41286e233 100644 --- a/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java +++ b/solr/cross-dc-manager/src/test/org/apache/solr/crossdc/manager/messageprocessor/TestMessageProcessor.java @@ -40,6 +40,7 @@ import org.apache.solr.crossdc.common.ResubmitBackoffPolicy; import org.apache.solr.crossdc.manager.consumer.ConsumerMetrics; import org.apache.solr.crossdc.manager.consumer.OtelMetrics; +import org.junit.After; import org.junit.Before; import org.junit.BeforeClass; import org.junit.Ignore; @@ -53,6 +54,7 @@ public class TestMessageProcessor { @Mock private CloudSolrClient solrClient; private SolrMessageProcessor processor; + private AutoCloseable mocks; private final ResubmitBackoffPolicy backoffPolicy = spy( @@ -70,13 +72,20 @@ public static void ensureWorkingMockito() { @Before public void setUp() { - MockitoAnnotations.openMocks(this); + mocks = MockitoAnnotations.openMocks(this); ConsumerMetrics metrics = Mockito.mock(OtelMetrics.class); processor = Mockito.spy(new SolrMessageProcessor(metrics, () -> solrClient, backoffPolicy)); Mockito.doNothing().when(processor).uncheckedSleep(anyLong()); } + @After + public void tearDown() throws Exception { + if (mocks != null) { + mocks.close(); + } + } + @Test public void testDocumentSanitization() { UpdateRequest request = spy(new UpdateRequest());