This is an automated email from the ASF dual-hosted git repository.

pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors.git


The following commit(s) were added to refs/heads/main by this push:
     new 3a06f7b70 solr 9 (#1813)
3a06f7b70 is described below

commit 3a06f7b709263c7c27f16c71a866c9ca780c30f6
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Aug 5 10:33:06 2026 +0100

    solr 9 (#1813)
    
    * solr 9
    
    * Update SolrTest.java
    
    * lru cache issue
---
 .scala-steward.conf                                |  4 +--
 project/Dependencies.scala                         |  4 +--
 .../connectors/solr/impl/SolrFlowStage.scala       | 30 +++++++++++++---------
 solr/src/test/java/docs/javadsl/SolrTest.java      | 14 +++++-----
 solr/src/test/resources/conf/solrconfig.xml        | 10 ++++----
 solr/src/test/scala/docs/scaladsl/SolrSpec.scala   | 18 +++++++------
 6 files changed, 44 insertions(+), 36 deletions(-)

diff --git a/.scala-steward.conf b/.scala-steward.conf
index 36b7a626f..d1ba2d8ce 100644
--- a/.scala-steward.conf
+++ b/.scala-steward.conf
@@ -1,8 +1,8 @@
 updates.pin  = [
   # pin to hadoop 3.4.x until 3.5.x becomes more widely adopted
   { groupId = "org.apache.hadoop", version = "3.4." }
-  # solrj 9.+ requires Java 11
-  { groupId = "org.apache.solr", version = "8." }
+  # stick with solr 9 until 10 is more widely adopted
+  { groupId = "org.apache.solr", version = "9." }
   # https://github.com/apache/pekko-connectors/issues/503
   { groupId = "com.couchbase.client", artifactId = "java-client", version = 
"2." }
   # activemq 6 is based on JakartaMS (only used in JMS tests and we test 
JakartaMS with Artemis)
diff --git a/project/Dependencies.scala b/project/Dependencies.scala
index c7ebe196a..4a2010042 100644
--- a/project/Dependencies.scala
+++ b/project/Dependencies.scala
@@ -503,8 +503,8 @@ object Dependencies {
         ExclusionRule("software.amazon.awssdk", "netty-nio-client")),
       "org.apache.pekko" %% "pekko-http" % PekkoHttpVersion) ++ Mockito)
 
-  val SolrjVersion = "8.11.4"
-  val SolrVersionForDocs = "8_11"
+  val SolrjVersion = "9.10.1"
+  val SolrVersionForDocs = "9_10"
 
   val Solr = Seq(
     libraryDependencies ++= Seq(
diff --git 
a/solr/src/main/scala/org/apache/pekko/stream/connectors/solr/impl/SolrFlowStage.scala
 
b/solr/src/main/scala/org/apache/pekko/stream/connectors/solr/impl/SolrFlowStage.scala
index 63e494e59..6a58f471e 100644
--- 
a/solr/src/main/scala/org/apache/pekko/stream/connectors/solr/impl/SolrFlowStage.scala
+++ 
b/solr/src/main/scala/org/apache/pekko/stream/connectors/solr/impl/SolrFlowStage.scala
@@ -112,19 +112,11 @@ private final class SolrFlowLogic[T, C](
       }
 
       message.routingFieldValue.foreach { routingFieldValue =>
-        val routingField = client match {
-          case csc: CloudSolrClient => {
-            val docCollection = 
Option(csc.getZkStateReader.getCollection(collection))
-            docCollection.flatMap { dc =>
-              Option(dc.getRouter.getRouteField(dc))
-            }
-          }
-          case _ => None
-        }
-        routingField.foreach { routingField =>
+        val routingField = getRoutingField()
+        routingField.foreach { rf =>
           message.idField.foreach { idField =>
-            if (routingField != idField)
-              doc.addField(routingField, routingFieldValue)
+            if (rf != idField)
+              doc.addField(rf, routingFieldValue)
           }
         }
       }
@@ -141,6 +133,20 @@ private final class SolrFlowLogic[T, C](
     client.add(collection, docs.asJava, settings.commitWithin)
   }
 
+  private def getRoutingField(): Option[String] = {
+    try {
+      client match {
+        case csc: CloudSolrClient =>
+          val provider = csc.getClusterStateProvider
+          val docCollection = provider.getCollection(collection)
+          Option(docCollection.getRouter.getRouteField(docCollection))
+        case _ => None
+      }
+    } catch {
+      case _: Exception => None
+    }
+  }
+
   private def deleteBulkToSolrByIds(messages: immutable.Seq[WriteMessage[T, 
C]]): UpdateResponse = {
     val docsIds = messages
       .filter { message =>
diff --git a/solr/src/test/java/docs/javadsl/SolrTest.java 
b/solr/src/test/java/docs/javadsl/SolrTest.java
index b84c4f42f..12f933321 100644
--- a/solr/src/test/java/docs/javadsl/SolrTest.java
+++ b/solr/src/test/java/docs/javadsl/SolrTest.java
@@ -44,9 +44,7 @@ import org.apache.pekko.testkit.javadsl.TestKit;
 import org.apache.solr.client.solrj.SolrClient;
 import org.apache.solr.client.solrj.SolrServerException;
 import org.apache.solr.client.solrj.beans.Field;
-import org.apache.solr.client.solrj.embedded.JettyConfig;
 import org.apache.solr.client.solrj.impl.CloudSolrClient;
-import org.apache.solr.client.solrj.impl.ZkClientClusterStateProvider;
 import org.apache.solr.client.solrj.io.SolrClientCache;
 import org.apache.solr.client.solrj.io.Tuple;
 import org.apache.solr.client.solrj.io.stream.CloudSolrStream;
@@ -58,6 +56,7 @@ import 
org.apache.solr.client.solrj.io.stream.expr.StreamFactory;
 import org.apache.solr.client.solrj.request.CollectionAdminRequest;
 import org.apache.solr.client.solrj.request.UpdateRequest;
 import org.apache.solr.client.solrj.response.UpdateResponse;
+import org.apache.solr.embedded.JettyConfig;
 import org.apache.solr.cloud.MiniSolrCloudCluster;
 import org.apache.solr.cloud.ZkTestServer;
 import org.apache.solr.common.SolrInputDocument;
@@ -391,6 +390,8 @@ public class SolrTest {
         List.of(0, 1, 2),
         CommittableOffsetBatch.committedOffsets.stream().map(o -> 
o.offset).toList());
 
+    solrClient.commit(collectionName);
+
     TupleStream stream = getTupleStream(collectionName);
 
     CompletionStage<List<String>> res2 =
@@ -846,7 +847,9 @@ public class SolrTest {
             testWorkingDir.toPath(),
             MiniSolrCloudCluster.DEFAULT_CLOUD_SOLR_XML,
             JettyConfig.builder().setContext("/solr").build(),
-            zkTestServer);
+            zkTestServer,
+            true);
+    cluster.uploadConfigSet(confDir.toPath(), "conf");
 
     // #init-client
 
@@ -855,10 +858,7 @@ public class SolrTest {
     // #init-client
     SolrTest.solrClient = solrClient;
 
-    ((ZkClientClusterStateProvider) solrClient.getClusterStateProvider())
-        .uploadConfig(confDir.toPath(), "conf");
-
-    
assertTrue(!solrClient.getZkStateReader().getClusterState().getLiveNodes().isEmpty());
+    assertTrue(!solrClient.getClusterStateProvider().getLiveNodes().isEmpty());
   }
 
   private static AtomicInteger number = new AtomicInteger(2);
diff --git a/solr/src/test/resources/conf/solrconfig.xml 
b/solr/src/test/resources/conf/solrconfig.xml
index 024aac0c7..700aeb093 100644
--- a/solr/src/test/resources/conf/solrconfig.xml
+++ b/solr/src/test/resources/conf/solrconfig.xml
@@ -26,24 +26,24 @@
     </updateHandler>
     <query>
         <maxBooleanClauses>1024</maxBooleanClauses>
-        <filterCache class="solr.FastLRUCache"
+        <filterCache class="org.apache.solr.search.FastLRUCache"
                      size="10"
                      initialSize="0"
                      autowarmCount="0"/>
-        <queryResultCache class="solr.LRUCache"
+        <queryResultCache class="org.apache.solr.search.LRUCache"
                           size="10"
                           initialSize="0"
                           autowarmCount="0"/>
-        <documentCache class="solr.LRUCache"
+        <documentCache class="org.apache.solr.search.LRUCache"
                        size="10"
                        initialSize="0"
                        autowarmCount="0"/>
         <cache name="perSegFilter"
-               class="solr.search.LRUCache"
+               class="org.apache.solr.search.LRUCache"
                size="10"
                initialSize="0"
                autowarmCount="10"
-               regenerator="solr.NoOpRegenerator"/>
+               regenerator="org.apache.solr.search.NoOpRegenerator"/>
 
         <enableLazyFieldLoading>true</enableLazyFieldLoading>
         <queryResultWindowSize>20</queryResultWindowSize>
diff --git a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala 
b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
index 27a75179b..4d9c77485 100644
--- a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
+++ b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
@@ -25,14 +25,14 @@ import pekko.stream.connectors.solr.scaladsl.{ SolrFlow, 
SolrSink, SolrSource }
 import pekko.stream.connectors.testkit.scaladsl.LogCapturing
 import pekko.stream.scaladsl.{ Sink, Source }
 import pekko.testkit.TestKit
-import org.apache.solr.client.solrj.embedded.JettyConfig
-import org.apache.solr.client.solrj.impl.{ CloudSolrClient, 
ZkClientClusterStateProvider }
+import org.apache.solr.client.solrj.impl.CloudSolrClient
 import org.apache.solr.client.solrj.io.stream.expr.{ StreamExpressionParser, 
StreamFactory }
 import org.apache.solr.client.solrj.io.stream.{ CloudSolrStream, 
StreamContext, TupleStream }
 import org.apache.solr.client.solrj.io.{ SolrClientCache, Tuple }
 import org.apache.solr.client.solrj.request.{ CollectionAdminRequest, 
UpdateRequest }
 import org.apache.solr.cloud.{ MiniSolrCloudCluster, ZkTestServer }
 import org.apache.solr.common.SolrInputDocument
+import org.apache.solr.embedded.JettyConfig
 import org.scalatest.concurrent.ScalaFutures
 import org.scalatest.BeforeAndAfterAll
 
@@ -323,6 +323,8 @@ class SolrSpec extends AnyWordSpec with Matchers with 
BeforeAndAfterAll with Sca
       // Make sure all messages was committed to kafka
       assert(List(0, 1, 2) == committedOffsets.map(_.offset))
 
+      solrClient.commit(collectionName)
+
       val stream = getTupleStream(collectionName)
 
       val res2 = SolrSource
@@ -495,6 +497,8 @@ class SolrSpec extends AnyWordSpec with Matchers with 
BeforeAndAfterAll with Sca
 
       deleteElements.futureValue
 
+      solrClient.commit(collectionName)
+
       val stream3 = getTupleStream(collectionName)
 
       val res2 = SolrSource
@@ -739,13 +743,11 @@ class SolrSpec extends AnyWordSpec with Matchers with 
BeforeAndAfterAll with Sca
       testWorkingDir.toPath,
       MiniSolrCloudCluster.DEFAULT_CLOUD_SOLR_XML,
       JettyConfig.builder.setContext("/solr").build,
-      zkTestServer)
-    solrClient.getClusterStateProvider
-      .asInstanceOf[ZkClientClusterStateProvider]
-      .uploadConfig(confDir.toPath, "conf")
-    solrClient.setIdField("router")
+      zkTestServer,
+      true)
+    cluster.uploadConfigSet(confDir.toPath, "conf")
 
-    assert(!solrClient.getZkStateReader.getClusterState.getLiveNodes.isEmpty)
+    assert(!solrClient.getClusterStateProvider.getLiveNodes.isEmpty)
   }
 
   private val number = new AtomicInteger(2)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to