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.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.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 void testKafkaExample() throws Exception { List.of(0, 1, 2), CommittableOffsetBatch.committedOffsets.stream().map(o -> o.offset).toList()); + solrClient.commit(collectionName); + TupleStream stream = getTupleStream(collectionName); CompletionStage> res2 = @@ -846,7 +847,9 @@ private static void setupCluster() throws Exception { 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 @@ private static void setupCluster() throws Exception { // #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 @@ 1024 - - - + regenerator="org.apache.solr.search.NoOpRegenerator"/> true 20 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)