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

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


The following commit(s) were added to refs/heads/main by this push:
     new 36bbf0fa6 slightly better support for Solr 9 client - (but full 
support won't be possible until we use Java 11 to build) (#481)
36bbf0fa6 is described below

commit 36bbf0fa68448b0efd86cac469454e956797c454
Author: PJ Fanning <[email protected]>
AuthorDate: Thu Feb 22 10:51:57 2024 +0100

    slightly better support for Solr 9 client - (but full support won't be 
possible until we use Java 11 to build) (#481)
    
    * temporarily drop support for getIdField
    
    * Update SolrFlowStage.scala
    
    * Update SolrFlowStage.scala
    
    * solr 9.4.1
    
    * use solr 9.5
    
    * fix some issues
    
    * comment out broken code
    
    * revert to solr8
    
    * Update .scala-steward.conf
    
    * review suggestion
---
 .scala-steward.conf                                          |  2 ++
 .../pekko/stream/connectors/solr/impl/SolrFlowStage.scala    | 11 +++++++----
 solr/src/test/java/docs/javadsl/SolrTest.java                | 12 +++++-------
 3 files changed, 14 insertions(+), 11 deletions(-)

diff --git a/.scala-steward.conf b/.scala-steward.conf
index 6a8d1dde8..4ef5d098f 100644
--- a/.scala-steward.conf
+++ b/.scala-steward.conf
@@ -5,6 +5,8 @@ updates.pin  = [
   { groupId = "org.springframework.boot", version = "2." }
   # spring-framework 6 requires Java 17
   { groupId = "org.springframework", version = "5." }
+  # solrj 9.+ requires Java 11
+  { groupId = "org.apache.solr", version = "8." }
   # mockito 5 requires Java 11 (only used in tests)
   { groupId = "org.mockito", version = "4." }
   # activemq 5.17+ requires Java 11 (only used in tests)
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 d03c87476..6a1e785b2 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
@@ -27,9 +27,8 @@ import org.apache.solr.client.solrj.response.UpdateResponse
 import org.apache.solr.common.SolrInputDocument
 
 import scala.annotation.tailrec
-import scala.util.control.NonFatal
-
 import scala.collection.immutable
+import scala.util.control.NonFatal
 
 /**
  * Internal API
@@ -114,8 +113,12 @@ private final class SolrFlowLogic[T, C](
 
       message.routingFieldValue.foreach { routingFieldValue =>
         val routingField = client match {
-          case csc: CloudSolrClient =>
-            Option(csc.getIdField)
+          case csc: CloudSolrClient => {
+            val docCollection = 
Option(csc.getZkStateReader.getCollection(collection))
+            docCollection.flatMap { dc =>
+              Option(dc.getRouter.getRouteField(dc))
+            }
+          }
           case _ => None
         }
         routingField.foreach { routingField =>
diff --git a/solr/src/test/java/docs/javadsl/SolrTest.java 
b/solr/src/test/java/docs/javadsl/SolrTest.java
index a853a0881..c4c6edb22 100644
--- a/solr/src/test/java/docs/javadsl/SolrTest.java
+++ b/solr/src/test/java/docs/javadsl/SolrTest.java
@@ -508,8 +508,8 @@ public class SolrTest {
         SolrSource.fromTupleStream(stream2)
             .map(
                 t -> {
-                  String id = t.fields.get("title").toString();
-                  String comment = t.fields.get("comment").toString();
+                  String id = t.getFields().get("title").toString();
+                  String comment = t.getFields().get("comment").toString();
                   Map<String, Object> m2 = new HashMap<>();
                   m2.put("set", (comment + " It's is a good book!!!"));
                   Map<String, Map<String, Object>> updates = new HashMap<>();
@@ -593,8 +593,8 @@ public class SolrTest {
         SolrSource.fromTupleStream(stream2)
             .map(
                 t -> {
-                  String id = t.fields.get("title").toString();
-                  String comment = t.fields.get("comment").toString();
+                  String id = t.getFields().get("title").toString();
+                  String comment = t.getFields().get("comment").toString();
                   Map<String, Object> m2 = new HashMap<>();
                   m2.put("set", (comment + " It's is a good book!!!"));
                   Map<String, Map<String, Object>> updates = new HashMap<>();
@@ -678,7 +678,7 @@ public class SolrTest {
         SolrSource.fromTupleStream(stream2)
             .map(
                 t -> {
-                  String id = t.fields.get("title").toString();
+                  String id = t.getFields().get("title").toString();
                   return 
WriteMessage.<SolrInputDocument>createDeleteByQueryMessage(
                       "title:\"" + id + "\"");
                 })
@@ -862,8 +862,6 @@ public class SolrTest {
     ((ZkClientClusterStateProvider) solrClient.getClusterStateProvider())
         .uploadConfig(confDir.toPath(), "conf");
 
-    solrClient.setIdField("router");
-
     
assertTrue(!solrClient.getZkStateReader().getClusterState().getLiveNodes().isEmpty());
   }
 


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

Reply via email to