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

juripetersen pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/wayang.git


The following commit(s) were added to refs/heads/main by this push:
     new 11d80e335 Add semantic filters with abstractions for semantic models 
as a platform (#808)
11d80e335 is described below

commit 11d80e33534808c70beceb88dfc55a8cc456e26d
Author: Juri Petersen <[email protected]>
AuthorDate: Mon Sep 14 14:03:46 2026 +0200

    Add semantic filters with abstractions for semantic models as a platform 
(#808)
    
    * implementation of semantic filters in Wayang
    
    * delete bin
    
    * delete unused file
    
    * spell check and grammar
    
    * Add missing license to guide
    
    * Add semantic operators with platforms per model
    
    * Update pom and fix address for Ollama
    
    * Mock Ollama API for test
    
    * Remove dotfiles
    
    * Remove default.properties for flink
    
    ---------
    
    Co-authored-by: mspruc <[email protected]>
---
 guides/semantic-operators-in-Wayang.md             | 156 +++++++
 pom.xml                                            |   2 +-
 .../scala/org/apache/wayang/api/DataQuanta.scala   |  33 ++
 .../org/apache/wayang/api/DataQuantaBuilder.scala  |  42 +-
 .../basic/operators/SemanticFilterOperator.java    |  77 ++++
 .../wayang/basic/operators/SemanticOperator.java   |  57 +++
 wayang-platforms/pom.xml                           |   1 +
 .../wayang/java/operators/JavaFilterOperator.java  |  11 +-
 wayang-platforms/wayang-semantic/pom.xml           |  75 ++++
 .../java/org/apache/wayang/semantic/Semantic.java  |  45 ++
 .../wayang/semantic/execution/OllamaExecutor.java  | 113 +++++
 .../operators/OllamaExecutionOperator.java         |  58 +++
 .../semantic/operators/OllamaFilterOperator.java   |  71 +++
 .../wayang/semantic/platform/OllamaPlatform.java   |  90 ++++
 .../wayang/semantic/plugin/SemanticPlugin.java     |  84 ++++
 .../wayang-semantic-ollama-defaults.properties     |  23 +
 .../org/apache/wayang/semantic/SemBenchTest.java   | 481 +++++++++++++++++++++
 17 files changed, 1415 insertions(+), 4 deletions(-)

diff --git a/guides/semantic-operators-in-Wayang.md 
b/guides/semantic-operators-in-Wayang.md
new file mode 100644
index 000000000..1a60faa77
--- /dev/null
+++ b/guides/semantic-operators-in-Wayang.md
@@ -0,0 +1,156 @@
+<!--
+
+  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.
+
+-->
+
+# Developing with Semantic Operators in Apache Wayang
+
+This guide explains how to define semantic operators, provide executable 
implementations, register multiple implementations, estimate their costs, and 
let Apache Wayang select an implementation during optimization.
+
+The example uses a semantic filter that classifies movie reviews as positive 
or negative through Ollama.
+
+## 1. Semantic operators
+
+A semantic operator is an operator much like any Wayang operator, however, it 
takes a prompt as input,
+that describes how it should act.
+
+For example:
+
+```java
+.semanticFilter(
+    "Analyze the review after the | and write either "
+        + "\"POSITIVE\" if the review has a positive sentiment and "
+        + "\"NEGATIVE\" if the review has a negative sentiment."
+)
+```
+
+The prompt describes the task. A `SemanticAlgorithm` provides one concrete 
implementation of that task.
+A `SemanticAlgorithm` could theoretically be any UDF you desire, there are no 
strict requirements on its implementation. 
+The implementation only that the requires the UDF is implemented as something 
that takes an input `Record` and a `prompt` and outputs
+whatever datatype is required by the operator.
+
+The Ollama local open-source model is used as an example for this guide, but 
may also be useful for quick development.
+
+We set up our local model hosting locally using Docker:
+
+```yaml
+ollama:
+    image: ollama/ollama:latest
+    container_name: apache-wayang-ollama
+    ports:
+        - "11434:11434"
+    volumes:
+        - ollama-data:/root/.ollama
+        - ./docker/ollama-init.sh:/ollama-init.sh
+    entrypoint: ["/bin/bash", "/ollama-init.sh"]
+    restart: always
+    tty: true
+    networks:
+        - wayang-network
+```
+
+```sh
+#!/bin/bash
+ollama serve &
+sleep 10
+ollama pull tinyllama
+wait
+```
+
+We setup the backend call to the model in Wayang:
+
+```java
+private static String callOllama(final String prompt) throws IOException, 
InterruptedException {
+    final String requestBody = String.format("{\"model\": \"%s\", \"prompt\": 
\"%s\", \"stream\": false}",
+            MODEL_NAME, escapeJson(prompt));
+
+    final HttpRequest request = 
HttpRequest.newBuilder().uri(URI.create(OLLAMA_API_URL))
+            .header("Content-Type", 
"application/json").POST(HttpRequest.BodyPublishers.ofString(requestBody))
+            .timeout(Duration.ofMinutes(2)).build();
+
+    final HttpResponse<String> response = httpClient.send(request, 
HttpResponse.BodyHandlers.ofString());
+
+    if (response.statusCode() != 200) {
+        throw new IOException("Ollama API error: " + response.body());
+    }
+
+    return parseOllamaResponse(response.body());
+}
+```
+
+and the semantic UDF:
+
+```java
+public static boolean isPositiveSentiment(final Review review, final String 
prompt)
+        throws IOException, InterruptedException {
+    final String response = callOllama(prompt + " | " + 
review.getReviewText());
+    return response.contains("POSITIVE");
+}
+```
+
+Now we can define our `SemanticAlgorithm`:
+
+```java
+final SemanticAlgorithm ollamaFilter = new SemanticAlgorithm();
+ollamaFilter.impl = (input, prompt) -> {
+    try {
+        return isPositiveSentiment((Review) input);
+    } catch (IOException | InterruptedException e) {
+        throw new RuntimeException("Ollama call failed", e);
+    }
+};
+```
+
+And also provide a UDF load estimator:
+
+```java
+ollamaFilter.loadProfileEstimator =
+    LoadProfileEstimators.createFromSpecification(
+            "wayang.semantic.ollama.model1.load",
+            configuration
+    );
+```
+
+Note that you should provide a separate configuration for each semantic 
operator.
+Now we register these UDFs with the `SemanticPlugin` this automatically 
constructs a new operator
+per UDF. Please note this may have implications for optimization time 
depending on your setup.
+
+```java
+final SemanticPlugin plugin = Semantic.plugin()
+        .withOperatorMapping(SemanticFilterOperator.class, ollamaFilter)
+        .withOperatorMapping(SemanticFilterOperator.class, ollamaFilter2)
+        .withOperatorMapping(SemanticFilterOperator.class, ollamaFilter3);
+
+final WayangContext wayangContext = new WayangContext()
+        .withPlugin(Java.basicPlugin())
+        .withPlugin(plugin);
+```
+
+Currently, we only have mappings for semantic operators in Java, so you also 
need the `Java.basicPlugin()`.
+Finally, you can construct your Wayang plan:
+
+```java
+final Collection<Long> positiveReviewCnt = 
planBuilder.loadCollection(loadReviews())
+        .filter(review -> "taken_3".equals(review.getId()))
+        .semanticFilter("Analyze the review after the | and write either 
\"POSITIVE\" if the review has a positive sentiment and \"NEGATIVE\" if the 
review has a negative sentiment.")
+            .withTargetModels(ollamaFilter, ollamaFilter2, ollamaFilter3)
+        .count()
+        .collect();
+```
+
+You need to provide semantic operators with their target models, even if you 
plan to use all models you've constructed.
+
diff --git a/pom.xml b/pom.xml
index 4d08cc103..f722d9921 100644
--- a/pom.xml
+++ b/pom.xml
@@ -23,7 +23,7 @@
     <parent>
         <groupId>org.apache</groupId>
         <artifactId>apache</artifactId>
-        <version>32</version>
+        <version>34</version>
     </parent>
 
     <groupId>org.apache.wayang</groupId>
diff --git 
a/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuanta.scala
 
b/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuanta.scala
index fb2dfeb77..f498d3eb4 100644
--- 
a/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuanta.scala
+++ 
b/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuanta.scala
@@ -44,6 +44,7 @@ import com.google.protobuf.ByteString
 import org.apache.wayang.api.python.function._
 import org.tensorflow.ndarray.NdArray
 
+import scala.collection.JavaConverters._
 import scala.collection.JavaConversions
 import scala.collection.JavaConversions._
 import scala.reflect._
@@ -633,6 +634,38 @@ class DataQuanta[Out: ClassTag](val operator: 
ElementaryOperator, outputIndex: I
     joinOperator
   }
 
+  def semanticFilterPrompt(prompt: String): DataQuanta[Out] = {
+    val dataSetType = org.apache.wayang.core.types.DataSetType.createDefault(
+      this.output.getType.getDataUnitType.toBasicDataUnitType
+    )
+
+    val filterOperator = new SemanticFilterOperator[Out](
+      dataSetType,
+      prompt
+    )
+
+    this.connectTo(filterOperator, 0)
+    wrap[Out](filterOperator)
+  }
+
+  def semanticFilterPrompt(prompt: String, targetModels: AnyRef * ): 
DataQuanta[Out] = {
+      val dataSetType = org.apache.wayang.core.types.DataSetType.createDefault(
+        this.output.getType.getDataUnitType.toBasicDataUnitType
+      )
+
+      val targetModelsSet: java.util.Set[Object] = 
targetModels.toSet.asJava.asInstanceOf[java.util.Set[Object]]
+      
+      val filterOperator = new SemanticFilterOperator[Out](
+        dataSetType,
+        prompt,
+        targetModelsSet
+      )
+
+      this.connectTo(filterOperator, 0)
+      wrap[Out](filterOperator)
+  }
+
+
   /**
     * Applies a spatial filter to this instance.
     *
diff --git 
a/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuantaBuilder.scala
 
b/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuantaBuilder.scala
index 9d37aa930..c399240fb 100644
--- 
a/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuantaBuilder.scala
+++ 
b/wayang-api/wayang-api-scala-java/src/main/scala/org/apache/wayang/api/DataQuantaBuilder.scala
@@ -45,10 +45,9 @@ import org.apache.iceberg.Schema
 import org.apache.iceberg.FileFormat
 import org.apache.iceberg.catalog.{Catalog, TableIdentifier}
 
-
-
 import scala.collection.mutable.ListBuffer
 import scala.reflect.ClassTag
+import scala.annotation.varargs
 
 /**
   * Trait/interface for builders of [[DataQuanta]]. The purpose of the 
builders is to provide a convenient
@@ -282,6 +281,9 @@ trait DataQuantaBuilder[+This <: DataQuantaBuilder[_, Out], 
Out] extends Logging
                          thatKeyUdf: SerializableFunction[ThatOut, Key]) =
     new JoinDataQuantaBuilder(this, that, thisKeyUdf, thatKeyUdf)
 
+  def semanticFilter(prompt: String) = 
+    new SemanticFilterDataQuantaBuilder[Out](this, prompt)
+
   /**
     * Feed the built [[DataQuanta]] into a spatial filter operator.
     * Requires the wayang-spatial plugin to be loaded.
@@ -2146,6 +2148,42 @@ class KeyedDataQuantaBuilder[Out, Key](private val 
dataQuantaBuilder: DataQuanta
 
 }
 
+/**
+  * [[DataQuantaBuilder]] implementation for 
[[org.apache.wayang.basic.operators.SemanticFilterOperator]]s.
+  *
+  * @param inputDataQuanta [[DataQuantaBuilder]] for the input [[DataQuanta]]
+  * @param udf             UDF for the [[SemanticFilterOperator]]
+  */
+class SemanticFilterDataQuantaBuilder[T](inputDataQuanta: DataQuantaBuilder[_, 
T], prompt: String)
+                                (implicit javaPlanBuilder: JavaPlanBuilder)
+  extends BasicDataQuantaBuilder[SemanticFilterDataQuantaBuilder[T], T] {
+
+    // Reuse the input TypeTrap to enforce type equality between input and 
output.
+    override def getOutputTypeTrap: TypeTrap = inputDataQuanta.outputTypeTrap
+
+    /** [[LoadProfileEstimator]] to estimate the [[LoadProfile]] of the 
[[udf]]. */
+    private var udfLoadProfileEstimator: LoadProfileEstimator = _
+
+    /** Selectivity of the filter predicate. */
+    private var selectivity: ProbabilisticDoubleInterval = _
+
+    /*
+
+    */
+    private val targetModels: ListBuffer[AnyRef] = ListBuffer()
+
+    @varargs def withTargetModels(models: AnyRef *): 
SemanticFilterDataQuantaBuilder[T] = {
+        models.foreach(targetModels.+=_)
+        this
+      }
+
+    override protected def build = applyTargetPlatforms(
+      inputDataQuanta.dataQuanta()
+          .semanticFilterPrompt(prompt, targetModels: _ *), 
+      this.getTargetPlatforms()
+    )
+  }
+
 class SpatialFilterDataQuantaBuilder[T](inputDataQuanta: DataQuantaBuilder[_, 
T],
                                         keySelector: SerializableFunction[T, _ 
<: SpatialGeometry],
                                         predicateType: SpatialPredicate,
diff --git 
a/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticFilterOperator.java
 
b/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticFilterOperator.java
new file mode 100644
index 000000000..f98898428
--- /dev/null
+++ 
b/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticFilterOperator.java
@@ -0,0 +1,77 @@
+/*
+ * 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.wayang.basic.operators;
+
+import java.util.Set;
+
+import org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator;
+import org.apache.wayang.core.types.DataSetType;
+
+public class SemanticFilterOperator<T> extends UnaryToUnaryOperator<T, T> 
implements SemanticOperator {
+
+    private final String prompt;
+
+    private final Set<Object> targetModels;
+
+    public SemanticFilterOperator(final DataSetType<T> type, final String 
prompt) {
+        super(type, type, false);
+        this.prompt = prompt;
+        this.targetModels = null;
+    }
+
+    public SemanticFilterOperator(final DataSetType<T> type, final String 
prompt, final Set<Object> targetModels) {
+        super(type, type, false);
+        this.prompt = prompt;
+        this.targetModels = targetModels;
+    }
+
+    public SemanticFilterOperator(final DataSetType<T> type) {
+        super(type, type, false);
+        this.prompt = "";
+        this.targetModels = null;
+    }
+
+    public SemanticFilterOperator(final DataSetType<T> inputType, final 
DataSetType<T> outputType,
+            final boolean isSupportingBroadcastInputs) {
+        super(inputType, outputType, isSupportingBroadcastInputs);
+        this.prompt = "";
+        this.targetModels = null;
+        }
+
+    public SemanticFilterOperator(final UnaryToUnaryOperator<T, T> that) {
+        super(that);
+        this.prompt = "";
+        this.targetModels = null;
+    }
+
+    @Override
+    public String getPrompt() {
+        return prompt;
+    }
+
+    @Override
+    public Set<Object> getTargetModels() {
+        return targetModels;
+    }
+
+    @Override
+    public void addTargetModel(final Object model) {
+        targetModels.add(model);
+    }
+}
diff --git 
a/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticOperator.java
 
b/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticOperator.java
new file mode 100644
index 000000000..29c00c80a
--- /dev/null
+++ 
b/wayang-commons/wayang-basic/src/main/java/org/apache/wayang/basic/operators/SemanticOperator.java
@@ -0,0 +1,57 @@
+/*
+ * 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.wayang.basic.operators;
+
+import java.util.Set;
+
+/**
+ * Common contract for operators whose UDF is not plain Java code but a 
natural-language
+ * {@code prompt} that gets executed by a large language model (LLM), e.g. a
+ * {@link SemanticFilterOperator} or a future {@code SemanticJoinOperator}.
+ * <p>
+ * This is deliberately kept as a slim interface rather than an abstract base 
class. Semantic
+ * operators still need to pick the operator shape that fits their 
arity/semantics, e.g. a
+ * filter extends {@link 
org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator} while a
+ * join would extend {@link 
org.apache.wayang.core.plan.wayangplan.BinaryToUnaryOperator}. As
+ * Java only allows single class inheritance, tying this abstraction to a 
specific operator
+ * base class would force every semantic operator into the same shape. 
Implementors are free
+ * to (and are expected to) extend whichever operator base class matches their 
actual
+ * input/output arity, while still sharing the prompt/target-model contract 
defined here.
+ */
+public interface SemanticOperator {
+
+    /**
+     * @return the prompt describing the semantic algorithm that this operator 
should execute
+     */
+    String getPrompt();
+
+    /**
+     * @return the set of models that are allowed to execute this operator's 
{@link #getPrompt()},
+     * or {@code null} if no target models have been configured
+     */
+    Set<Object> getTargetModels();
+
+    /**
+     * Adds a model to the set of models that are allowed to execute this 
operator's
+     * {@link #getPrompt()}.
+     *
+     * @param model the model to add
+     */
+    void addTargetModel(Object model);
+}
diff --git a/wayang-platforms/pom.xml b/wayang-platforms/pom.xml
index a4b062563..ece339f70 100644
--- a/wayang-platforms/pom.xml
+++ b/wayang-platforms/pom.xml
@@ -47,6 +47,7 @@
         <module>wayang-presto</module>
         <module>wayang-trino</module>
         <module>wayang-tensorflow</module>
+        <module>wayang-semantic</module>
     </modules>
 
     <dependencyManagement>
diff --git 
a/wayang-platforms/wayang-java/src/main/java/org/apache/wayang/java/operators/JavaFilterOperator.java
 
b/wayang-platforms/wayang-java/src/main/java/org/apache/wayang/java/operators/JavaFilterOperator.java
index 757220e34..a11e7625d 100644
--- 
a/wayang-platforms/wayang-java/src/main/java/org/apache/wayang/java/operators/JavaFilterOperator.java
+++ 
b/wayang-platforms/wayang-java/src/main/java/org/apache/wayang/java/operators/JavaFilterOperator.java
@@ -50,6 +50,16 @@ public class JavaFilterOperator<Type>
         implements JavaExecutionOperator {
 
 
+    
+    /**
+     * Creates a new instance.
+     *
+     * @param type type of the dataset elements
+     */
+    public JavaFilterOperator(PredicateDescriptor<Type> predicateDescriptor) {
+        super(predicateDescriptor);
+    }
+
     /**
      * Creates a new instance.
      *
@@ -119,5 +129,4 @@ public class JavaFilterOperator<Type>
         assert index <= this.getNumOutputs() || (index == 0 && 
this.getNumOutputs() == 0);
         return Collections.singletonList(StreamChannel.DESCRIPTOR);
     }
-
 }
diff --git a/wayang-platforms/wayang-semantic/pom.xml 
b/wayang-platforms/wayang-semantic/pom.xml
new file mode 100644
index 000000000..e32f80676
--- /dev/null
+++ b/wayang-platforms/wayang-semantic/pom.xml
@@ -0,0 +1,75 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+  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.
+  -->
+<project xmlns="http://maven.apache.org/POM/4.0.0"; 
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"; 
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <modelVersion>4.0.0</modelVersion>
+
+    <parent>
+        <artifactId>wayang-platforms</artifactId>
+        <groupId>org.apache.wayang</groupId>
+        <version>1.1.2-SNAPSHOT</version>
+    </parent>
+
+    <artifactId>wayang-semantic</artifactId>
+
+    <name>Wayang Platform Semantic</name>
+    <description>
+        Wayang implementation of semantic operators
+    </description>
+
+    <properties>
+        
<java-module-name>org.apache.wayang.platform.semantic</java-module-name>
+    </properties>
+
+    <dependencies>
+        <dependency>
+            <groupId>org.apache.wayang</groupId>
+            <artifactId>wayang-java</artifactId>
+            <version>1.1.2-SNAPSHOT</version>
+        </dependency>
+        <dependency>
+            <groupId>org.apache.wayang</groupId>
+            <artifactId>wayang-basic</artifactId>
+            <version>1.1.2-SNAPSHOT</version>
+        </dependency>
+        <dependency>
+            <groupId>org.apache.wayang</groupId>
+            <artifactId>wayang-java</artifactId>
+            <version>1.1.2-SNAPSHOT</version>
+        </dependency>
+        <dependency>
+            <groupId>org.junit.jupiter</groupId>
+            <artifactId>junit-jupiter</artifactId>
+            <version>5.10.2</version>
+            <scope>test</scope>
+        </dependency>
+        <dependency>
+            <groupId>com.fasterxml.jackson.core</groupId>
+            <artifactId>jackson-core</artifactId>
+            <version>2.16.1</version> 
+            <scope>test</scope>
+        </dependency>
+        <dependency>
+            <groupId>org.apache.wayang</groupId>
+            <artifactId>wayang-api-scala-java</artifactId>
+            <version>1.1.2-SNAPSHOT</version>
+            <scope>test</scope>
+        </dependency>
+    </dependencies>
+</project>
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/Semantic.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/Semantic.java
new file mode 100644
index 000000000..83f497008
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/Semantic.java
@@ -0,0 +1,45 @@
+/*
+ * 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.wayang.semantic;
+
+import org.apache.wayang.semantic.plugin.SemanticPlugin;
+import org.apache.wayang.java.platform.JavaPlatform;
+
+public class Semantic {
+    private final static SemanticPlugin PLUGIN = new SemanticPlugin();
+
+    /**
+     * Retrieve the {@link SemanticPlugin}.
+     *
+     * @return the {@link SemanticPlugin}
+     */
+    public static SemanticPlugin plugin() {
+        return PLUGIN;
+    }
+
+
+    /**
+     * Retrieve the {@link SemanticPlatform}.
+     *
+     * @return the {@link SemanticPlatform}
+     */
+    public static JavaPlatform platform() {
+        return JavaPlatform.getInstance();
+    }
+}
\ No newline at end of file
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/execution/OllamaExecutor.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/execution/OllamaExecutor.java
new file mode 100644
index 000000000..706d90241
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/execution/OllamaExecutor.java
@@ -0,0 +1,113 @@
+/*
+ * 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.wayang.semantic.execution;
+
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.List;
+
+import org.apache.wayang.core.api.Job;
+import org.apache.wayang.core.api.exception.WayangException;
+import org.apache.wayang.core.optimizer.OptimizationContext;
+import org.apache.wayang.core.plan.executionplan.ExecutionTask;
+import org.apache.wayang.core.plan.wayangplan.ExecutionOperator;
+import org.apache.wayang.core.platform.ChannelInstance;
+import org.apache.wayang.core.platform.Executor;
+import org.apache.wayang.core.platform.PartialExecution;
+import org.apache.wayang.core.platform.PushExecutorTemplate;
+import org.apache.wayang.core.platform.lineage.ExecutionLineageNode;
+import org.apache.wayang.core.util.Formats;
+import org.apache.wayang.core.util.Tuple;
+import org.apache.wayang.semantic.operators.OllamaExecutionOperator;
+import org.apache.wayang.semantic.platform.OllamaPlatform;
+
+/**
+ * {@link Executor} implementation for the {@link OllamaPlatform}.
+ */
+public class OllamaExecutor extends PushExecutorTemplate {
+
+    private final OllamaPlatform platform;
+
+    public OllamaExecutor(final OllamaPlatform platform, final Job job) {
+        super(job);
+        this.platform = platform;
+    }
+
+    @Override
+    public OllamaPlatform getPlatform() {
+        return this.platform;
+    }
+
+    @Override
+    protected Tuple<List<ChannelInstance>, PartialExecution> execute(
+            final ExecutionTask task,
+            final List<ChannelInstance> inputChannelInstances,
+            final OptimizationContext.OperatorContext producerOperatorContext,
+            final boolean isRequestEagerExecution
+    ) {
+        final ChannelInstance[] outputChannelInstances = 
task.getOperator().createOutputChannelInstances(
+                this, task, producerOperatorContext, inputChannelInstances
+        );
+
+        final Collection<ExecutionLineageNode> executionLineageNodes;
+        final Collection<ChannelInstance> producedChannelInstances;
+        this.job.reportProgress(task.getOperator().getName(), 50);
+        final long startTime = System.currentTimeMillis();
+        try {
+            final Tuple<Collection<ExecutionLineageNode>, 
Collection<ChannelInstance>> results =
+                    cast(task.getOperator()).evaluate(
+                            toArray(inputChannelInstances),
+                            outputChannelInstances,
+                            this,
+                            producerOperatorContext
+                    );
+            executionLineageNodes = results.getField0();
+            producedChannelInstances = results.getField1();
+        } catch (Exception e) {
+            throw new WayangException(String.format("Executing %s failed.", 
task), e);
+        }
+        final long endTime = System.currentTimeMillis();
+        final long executionDuration = endTime - startTime;
+
+        this.job.reportProgress(task.getOperator().getName(), 100);
+
+        final PartialExecution partialExecution = 
this.createPartialExecution(executionLineageNodes, executionDuration);
+
+        if (partialExecution == null && executionDuration > 10) {
+            this.logger.warn("Execution of {} took suspiciously long ({}).", 
task, Formats.formatDuration(executionDuration));
+        }
+
+        this.registerMeasuredCardinalities(producedChannelInstances);
+
+        if (isRequestEagerExecution && partialExecution == null) {
+            this.logger.info("{} was not executed eagerly as requested.", 
task);
+        }
+
+        return new Tuple<>(Arrays.asList(outputChannelInstances), 
partialExecution);
+    }
+
+    private static OllamaExecutionOperator cast(final ExecutionOperator 
executionOperator) {
+        return (OllamaExecutionOperator) executionOperator;
+    }
+
+    private static ChannelInstance[] toArray(final List<ChannelInstance> 
channelInstances) {
+        final ChannelInstance[] array = new 
ChannelInstance[channelInstances.size()];
+        return channelInstances.toArray(array);
+    }
+}
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaExecutionOperator.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaExecutionOperator.java
new file mode 100644
index 000000000..b635a95b4
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaExecutionOperator.java
@@ -0,0 +1,58 @@
+/*
+ * 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.wayang.semantic.operators;
+
+import java.util.Collection;
+
+import org.apache.wayang.core.optimizer.OptimizationContext;
+import org.apache.wayang.core.plan.wayangplan.ExecutionOperator;
+import org.apache.wayang.core.platform.ChannelInstance;
+import org.apache.wayang.core.platform.lineage.ExecutionLineageNode;
+import org.apache.wayang.core.util.Tuple;
+import org.apache.wayang.semantic.execution.OllamaExecutor;
+import org.apache.wayang.semantic.platform.OllamaPlatform;
+
+/**
+ * Execution operator for the {@link OllamaPlatform}.
+ */
+public interface OllamaExecutionOperator extends ExecutionOperator {
+
+    @Override
+    default OllamaPlatform getPlatform() {
+        return OllamaPlatform.getInstance();
+    }
+
+    /**
+     * Evaluates this operator by calling out to Ollama. Mirrors
+     * {@code 
org.apache.wayang.java.operators.JavaExecutionOperator#evaluate}, just for the
+     * {@link OllamaPlatform} instead of the Java platform.
+     *
+     * @param inputs          {@link ChannelInstance}s that satisfy the inputs 
of this operator
+     * @param outputs         {@link ChannelInstance}s that collect the 
outputs of this operator
+     * @param ollamaExecutor  that executes this instance
+     * @param operatorContext optimization information for this instance
+     * @return {@link Collection}s of what has been executed and produced
+     */
+    Tuple<Collection<ExecutionLineageNode>, Collection<ChannelInstance>> 
evaluate(
+            ChannelInstance[] inputs,
+            ChannelInstance[] outputs,
+            OllamaExecutor ollamaExecutor,
+            OptimizationContext.OperatorContext operatorContext);
+
+}
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaFilterOperator.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaFilterOperator.java
new file mode 100644
index 000000000..40d9ee7ae
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/operators/OllamaFilterOperator.java
@@ -0,0 +1,71 @@
+/*
+ * 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.wayang.semantic.operators;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+import org.apache.wayang.basic.operators.SemanticFilterOperator;
+import org.apache.wayang.core.plan.wayangplan.ExecutionOperator;
+import org.apache.wayang.core.platform.ChannelDescriptor;
+import org.apache.wayang.core.types.DataSetType;
+import org.apache.wayang.java.channels.CollectionChannel;
+import org.apache.wayang.java.channels.StreamChannel;
+
+/**
+ * Base for {@link OllamaExecutionOperator}s that implement a {@link 
SemanticFilterOperator} by calling
+ * out to one specific model. Every concrete subclass hardcodes exactly one 
model's evaluation logic
+ * instead of taking it as an injected UDF, so {@link #getPrompt()} (inherited 
from the matched logical
+ * {@link SemanticFilterOperator}) remains the single source of truth for the 
prompt; a subclass is free
+ * to use it or to ignore it in favor of its own hardcoded behavior.
+ * <p>
+ * Input/output channels are the same {@link StreamChannel}/{@link 
CollectionChannel} that
+ * {@code JavaFilterOperator} supports, so data can flow between the 
surrounding Java pipeline and an
+ * Ollama-backed operator without any channel conversion.
+ */
+public abstract class OllamaFilterOperator<Type> extends 
SemanticFilterOperator<Type> implements OllamaExecutionOperator {
+
+    protected OllamaFilterOperator(final DataSetType<Type> type, final String 
prompt) {
+        super(type, prompt);
+    }
+
+    /**
+     * Creates a new instance of the same concrete subclass for the given 
type/prompt.
+     */
+    protected abstract OllamaFilterOperator<Type> 
newInstance(DataSetType<Type> type, String prompt);
+
+    @Override
+    protected final ExecutionOperator createCopy() {
+        return this.newInstance(this.getInputType(), this.getPrompt());
+    }
+
+    @Override
+    public List<ChannelDescriptor> getSupportedInputChannels(final int index) {
+        assert index <= this.getNumInputs() || (index == 0 && 
this.getNumInputs() == 0);
+        if (this.getInput(index).isBroadcast()) return 
Collections.singletonList(CollectionChannel.DESCRIPTOR);
+        return Arrays.asList(CollectionChannel.DESCRIPTOR, 
StreamChannel.DESCRIPTOR);
+    }
+
+    @Override
+    public List<ChannelDescriptor> getSupportedOutputChannels(final int index) 
{
+        assert index <= this.getNumOutputs() || (index == 0 && 
this.getNumOutputs() == 0);
+        return Collections.singletonList(StreamChannel.DESCRIPTOR);
+    }
+}
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/platform/OllamaPlatform.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/platform/OllamaPlatform.java
new file mode 100644
index 000000000..1cad4d6b1
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/platform/OllamaPlatform.java
@@ -0,0 +1,90 @@
+/*
+ * 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.wayang.semantic.platform;
+
+import org.apache.wayang.core.api.Configuration;
+import org.apache.wayang.core.optimizer.costs.LoadProfileToTimeConverter;
+import org.apache.wayang.core.optimizer.costs.LoadToTimeConverter;
+import org.apache.wayang.core.optimizer.costs.TimeToCostConverter;
+import org.apache.wayang.core.platform.Executor;
+import org.apache.wayang.core.platform.Platform;
+import org.apache.wayang.core.util.ReflectionUtils;
+import org.apache.wayang.semantic.execution.OllamaExecutor;
+
+/**
+ * {@link Platform} for operators that are executed by calling out to an 
Ollama-hosted LLM.
+ * <p>
+ * Kept as its own {@link Platform} (rather than piggy-backing on
+ * {@link org.apache.wayang.java.platform.JavaPlatform}) so that the cost of 
an LLM call is modeled and
+ * optimized on its own terms, instead of being reported as if it were free 
Java-Stream work.
+ */
+public class OllamaPlatform extends Platform {
+
+    private static final String PLATFORM_NAME = "Ollama";
+
+    private static final String CONFIG_NAME = "semantic.ollama";
+
+    private static final String DEFAULT_CONFIG_FILE = 
"wayang-semantic-ollama-defaults.properties";
+
+    private static OllamaPlatform instance = null;
+
+    public static OllamaPlatform getInstance() {
+        if (instance == null) {
+            instance = new OllamaPlatform();
+        }
+        return instance;
+    }
+
+    private OllamaPlatform() {
+        super(PLATFORM_NAME, CONFIG_NAME);
+    }
+
+    @Override
+    public void configureDefaults(final Configuration configuration) {
+        configuration.load(ReflectionUtils.loadResource(DEFAULT_CONFIG_FILE));
+    }
+
+    @Override
+    public Executor.Factory getExecutorFactory() {
+        return job -> new OllamaExecutor(this, job);
+    }
+
+    @Override
+    public LoadProfileToTimeConverter createLoadProfileToTimeConverter(final 
Configuration configuration) {
+        final int cpuMhz = (int) 
configuration.getLongProperty("wayang.semantic.ollama.cpu.mhz");
+        final int numCores = (int) 
configuration.getLongProperty("wayang.semantic.ollama.cores");
+        final double hdfsMsPerMb = 
configuration.getDoubleProperty("wayang.semantic.ollama.hdfs.ms-per-mb");
+        final double stretch = 
configuration.getDoubleProperty("wayang.semantic.ollama.stretch");
+        return LoadProfileToTimeConverter.createTopLevelStretching(
+                LoadToTimeConverter.createLinearCoverter(1 / (numCores * 
cpuMhz * 1000d)),
+                LoadToTimeConverter.createLinearCoverter(hdfsMsPerMb / 
1000000d),
+                LoadToTimeConverter.createLinearCoverter(0),
+                (cpuEstimate, diskEstimate, networkEstimate) -> 
cpuEstimate.plus(diskEstimate).plus(networkEstimate),
+                stretch
+        );
+    }
+
+    @Override
+    public TimeToCostConverter createTimeToCostConverter(final Configuration 
configuration) {
+        return new TimeToCostConverter(
+                
configuration.getDoubleProperty("wayang.semantic.ollama.costs.fix"),
+                
configuration.getDoubleProperty("wayang.semantic.ollama.costs.per-ms")
+        );
+    }
+}
diff --git 
a/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/plugin/SemanticPlugin.java
 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/plugin/SemanticPlugin.java
new file mode 100644
index 000000000..2b0964f55
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/java/org/apache/wayang/semantic/plugin/SemanticPlugin.java
@@ -0,0 +1,84 @@
+/*
+ * 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.wayang.semantic.plugin;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Set;
+
+import org.apache.wayang.core.api.Configuration;
+import org.apache.wayang.core.mapping.Mapping;
+import org.apache.wayang.core.mapping.PlanTransformation;
+import org.apache.wayang.core.optimizer.channels.ChannelConversion;
+import org.apache.wayang.core.platform.Platform;
+import org.apache.wayang.core.plugin.Plugin;
+import org.apache.wayang.java.Java;
+import org.apache.wayang.java.platform.JavaPlatform;
+
+public class SemanticPlugin implements Plugin {
+    private final List<Mapping> mappings;
+
+    private SemanticPlugin(final List<Mapping> mappings) {
+        this.mappings = mappings;
+    }
+
+    public SemanticPlugin() {
+        this.mappings = List.of();
+    }
+
+    @Override
+    public Collection<Mapping> getMappings() {
+        return mappings;
+    }
+
+    @Override
+    public Collection<Platform> getRequiredPlatforms() {
+        final Set<Platform> platforms = new LinkedHashSet<>();
+        // The surrounding pipeline (sources, sinks, plain transformations) 
always runs on Java.
+        platforms.add(JavaPlatform.getInstance());
+        for (final Mapping mapping : this.mappings) {
+            for (final PlanTransformation transformation : 
mapping.getTransformations()) {
+                platforms.addAll(transformation.getTargetPlatforms());
+            }
+        }
+        return platforms;
+    }
+
+    @Override
+    public Collection<ChannelConversion> getChannelConversions() {
+        return Java.basicPlugin().getChannelConversions();
+    }
+
+    @Override
+    public void setProperties(final Configuration configuration) {
+    }
+
+    /**
+     * Registers a {@link Mapping} that rewrites a semantic operator (e.g. a
+     * {@link org.apache.wayang.basic.operators.SemanticFilterOperator}) into 
a physical operator for one
+     * specific model implementation.
+     */
+    public SemanticPlugin withMapping(final Mapping mapping) {
+        final List<Mapping> nextMappings = new ArrayList<>(this.mappings);
+        nextMappings.add(mapping);
+        return new SemanticPlugin(nextMappings);
+    }
+}
diff --git 
a/wayang-platforms/wayang-semantic/src/main/resources/wayang-semantic-ollama-defaults.properties
 
b/wayang-platforms/wayang-semantic/src/main/resources/wayang-semantic-ollama-defaults.properties
new file mode 100644
index 000000000..13106b24e
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/main/resources/wayang-semantic-ollama-defaults.properties
@@ -0,0 +1,23 @@
+#
+# 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.
+#
+
+wayang.semantic.ollama.cpu.mhz = 2700
+wayang.semantic.ollama.cores = 1
+wayang.semantic.ollama.hdfs.ms-per-mb = 2.7
+wayang.semantic.ollama.stretch = 1
+wayang.semantic.ollama.costs.fix = 0.0
+wayang.semantic.ollama.costs.per-ms = 1.0
diff --git 
a/wayang-platforms/wayang-semantic/src/test/java/org/apache/wayang/semantic/SemBenchTest.java
 
b/wayang-platforms/wayang-semantic/src/test/java/org/apache/wayang/semantic/SemBenchTest.java
new file mode 100644
index 000000000..30f35ccf2
--- /dev/null
+++ 
b/wayang-platforms/wayang-semantic/src/test/java/org/apache/wayang/semantic/SemBenchTest.java
@@ -0,0 +1,481 @@
+/*
+ * 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.wayang.semantic;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.net.InetSocketAddress;
+import java.net.URI;
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
+import java.util.stream.Stream;
+
+import org.apache.wayang.api.JavaPlanBuilder;
+import org.apache.wayang.basic.operators.SemanticFilterOperator;
+import org.apache.wayang.core.api.Configuration;
+import org.apache.wayang.core.api.WayangContext;
+import org.apache.wayang.core.mapping.Mapping;
+import org.apache.wayang.core.mapping.OperatorPattern;
+import org.apache.wayang.core.mapping.PlanTransformation;
+import org.apache.wayang.core.mapping.ReplacementSubplanFactory;
+import org.apache.wayang.core.mapping.SubplanPattern;
+import org.apache.wayang.core.optimizer.OptimizationContext;
+import org.apache.wayang.core.plan.wayangplan.ExecutionOperator;
+import org.apache.wayang.core.platform.ChannelInstance;
+import org.apache.wayang.core.platform.lineage.ExecutionLineageNode;
+import org.apache.wayang.core.types.DataSetType;
+import org.apache.wayang.core.util.Tuple;
+import org.apache.wayang.java.Java;
+import org.apache.wayang.java.channels.JavaChannelInstance;
+import org.apache.wayang.java.channels.StreamChannel;
+import org.apache.wayang.semantic.execution.OllamaExecutor;
+import org.apache.wayang.semantic.operators.OllamaFilterOperator;
+import org.apache.wayang.semantic.platform.OllamaPlatform;
+import org.apache.wayang.semantic.plugin.SemanticPlugin;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.sun.net.httpserver.HttpServer;
+
+class SemBenchTest {
+
+    /**
+     * In-process stand-in for a real Ollama server: it never leaves the JVM 
and never needs a model
+     * to be pulled, so the test runs unmodified in CI. It answers with the 
same response envelope a
+     * real Ollama server would ("{"response":"..."}"), classifying sentiment 
by a keyword heuristic
+     * on the prompt text, so the assertions below exercise genuine end-to-end 
behavior of the plan
+     * (filtering, the semantic operator, and JSON parsing) rather than a 
canned constant.
+     */
+    private static final String[] NEGATIVE_KEYWORDS = {"terrible", 
"disappointed", "boring", "not recommend"};
+
+    private HttpServer mockOllamaServer;
+
+    @BeforeEach
+    void startMockOllamaServer() throws IOException {
+        this.mockOllamaServer = HttpServer.create(new 
InetSocketAddress("127.0.0.1", 0), 0);
+        this.mockOllamaServer.createContext("/api/generate", exchange -> {
+            final String requestBody = new 
String(exchange.getRequestBody().readAllBytes(), StandardCharsets.UTF_8);
+            final boolean isNegative = Arrays.stream(NEGATIVE_KEYWORDS)
+                    .anyMatch(keyword -> 
requestBody.toLowerCase().contains(keyword));
+            final String responseBody = String.format("{\"response\":\"%s\"}", 
isNegative ? "NEGATIVE" : "POSITIVE");
+            final byte[] responseBytes = 
responseBody.getBytes(StandardCharsets.UTF_8);
+
+            exchange.getResponseHeaders().add("Content-Type", 
"application/json");
+            exchange.sendResponseHeaders(200, responseBytes.length);
+            try (OutputStream responseStream = exchange.getResponseBody()) {
+                responseStream.write(responseBytes);
+            }
+        });
+        this.mockOllamaServer.start();
+
+        final int port = this.mockOllamaServer.getAddress().getPort();
+        OllamaSemanticFilter.setApiUrl("http://127.0.0.1:"; + port + 
"/api/generate");
+    }
+
+    @AfterEach
+    void stopMockOllamaServer() {
+        this.mockOllamaServer.stop(0);
+    }
+
+    private static List<Review> loadReviews() {
+        return Arrays.asList(new Review("taken_1", "The movie was fantastic. 
Great acting and an engaging story."),
+                new Review("taken_2", "I was disappointed. The plot was boring 
and too long."),
+                new Review("taken_3", "Absolutely loved it! One of the best 
movies I have seen this year."),
+                new Review("taken_3", "Terrible experience. I would not 
recommend it to anyone."),
+                new Review("taken_4", "It was okay. Not great, not 
terrible."));
+    }
+
+    @Test
+    void testSemBenchMoviesWithOllama() {
+        final Configuration configuration = new Configuration();
+        configuration.setProperty("wayang.java.filter.load", """
+            {
+            "in":1,
+            "out":1,
+            "cpu":"${25*in0 + 350000}",
+            "ram":"100000",
+            "p":0.9
+            }
+            """
+        );
+        configuration.setProperty("wayang.semantic.ollama.model1.load",
+            """
+            {
+            "in": 1,
+            "out": 1,
+            "cpu": "${500*in0 + 56789}",
+            "ram": "10000",
+            "disk": "0",
+            "net": "0",
+            "p": 0.9,
+            "overhead": 0,
+            "ru": "${wayang:logGrowth(0.1, 0.1, 1000000, in0)}"
+            }
+            """);
+        configuration.setProperty("wayang.semantic.ollama.model2.load",
+            """
+            {
+            "in": 1,
+            "out": 1,
+            "cpu": "${500*in0 + 56789}",
+            "ram": "10000",
+            "disk": "0",
+            "net": "0",
+            "p": 0.9,
+            "overhead": 0,
+            "ru": "${wayang:logGrowth(0.1, 0.1, 1000000, in0)}"
+            }
+            """);
+        configuration.setProperty("wayang.semantic.ollama.model3.load",
+            """
+            {
+            "in": 1,
+            "out": 1,
+            "cpu": "${50*in0 + 5678}",
+            "ram": "1000",
+            "disk": "0",
+            "net": "0",
+            "p": 0.9,
+            "overhead": 0,
+            "ru": "${wayang:logGrowth(0.1, 0.1, 1000000, in0)}"
+            }
+            """
+        );
+
+        // Each Ollama model is its own physical operator + Mapping; the 
optimizer picks among them by cost.
+        // There is exactly one prompt, owned by the logical 
SemanticFilterOperator (set via .semanticFilter(...)),
+        // and each physical operator decides for itself whether it needs that 
prompt.
+        final SemanticPlugin plugin = Semantic.plugin()
+                .withMapping(new OllamaModel1FilterMapping())
+                .withMapping(new OllamaModel2FilterMapping())
+                .withMapping(new OllamaModel3FilterMapping());
+
+        final WayangContext wayangContext = new WayangContext(configuration)
+                .withPlugin(Java.basicPlugin())
+                .withPlugin(plugin);
+        final JavaPlanBuilder planBuilder = new JavaPlanBuilder(wayangContext);
+
+        final Collection<Long> positiveReviewCnt = 
planBuilder.loadCollection(loadReviews())
+                .filter(review -> "taken_3".equals(review.getId()))
+                .semanticFilter("Analyze the review after the | and write 
either \"POSITIVE\" if the review has a positive sentiment and \"NEGATIVE\" if 
the review has a negative sentiment.")
+                    .withTargetModels(OllamaModel1FilterOperator.class, 
OllamaModel2FilterOperator.class, OllamaModel3FilterOperator.class)
+                .count()
+                .collect();
+
+        // Of the two "taken_3" reviews, only the "Absolutely loved it!" one 
is positive; the mock
+        // Ollama server classifies by keyword, so this checks the plan 
actually ran the semantic
+        // filter rather than merely completing without error.
+        assertEquals(1, positiveReviewCnt.size());
+        assertEquals(1L, positiveReviewCnt.iterator().next());
+    }
+}
+
+class Review {
+    private final String id;
+    private final String reviewText;
+
+    public Review(final String id, final String reviewText) {
+        this.id = id;
+        this.reviewText = reviewText;
+    }
+
+    public String getId() {
+        return id;
+    }
+
+    public String getReviewText() {
+        return reviewText;
+    }
+
+    @Override
+    public String toString() {
+        return "Review{id='" + id + "', reviewText='" + reviewText + "'}";
+    }
+}
+
+/**
+ * Physical operator for the "model1" Ollama implementation: a fixed sentiment 
prompt, ignores the
+ * logical operator's {@link #getPrompt()}.
+ */
+final class OllamaModel1FilterOperator<Type> extends 
OllamaFilterOperator<Type> {
+
+    OllamaModel1FilterOperator(final DataSetType<Type> type, final String 
prompt) {
+        super(type, prompt);
+    }
+
+    @Override
+    @SuppressWarnings("unchecked")
+    public Tuple<Collection<ExecutionLineageNode>, 
Collection<ChannelInstance>> evaluate(
+            final ChannelInstance[] inputs,
+            final ChannelInstance[] outputs,
+            final OllamaExecutor ollamaExecutor,
+            final OptimizationContext.OperatorContext operatorContext) {
+        assert inputs.length == this.getNumInputs();
+        assert outputs.length == this.getNumOutputs();
+
+        final Stream<Review> filtered = ((JavaChannelInstance) 
inputs[0]).<Review>provideStream().filter(review -> {
+            try {
+                return OllamaSemanticFilter.isPositiveSentiment(review);
+            } catch (IOException | InterruptedException e) {
+                throw new RuntimeException("Ollama call failed", e);
+            }
+        });
+        ((StreamChannel.Instance) outputs[0]).accept(filtered);
+
+        return ExecutionOperator.modelLazyExecution(inputs, outputs, 
operatorContext);
+    }
+
+    @Override
+    public String getLoadProfileEstimatorConfigurationKey() {
+        return "wayang.semantic.ollama.model1.load";
+    }
+
+    @Override
+    protected OllamaFilterOperator<Type> newInstance(final DataSetType<Type> 
type, final String prompt) {
+        return new OllamaModel1FilterOperator<>(type, prompt);
+    }
+}
+
+/**
+ * Physical operator for the "model2" Ollama implementation: a different fixed 
sentiment prompt, also
+ * ignores the logical operator's {@link #getPrompt()}.
+ */
+final class OllamaModel2FilterOperator<Type> extends 
OllamaFilterOperator<Type> {
+
+    OllamaModel2FilterOperator(final DataSetType<Type> type, final String 
prompt) {
+        super(type, prompt);
+    }
+
+    @Override
+    @SuppressWarnings("unchecked")
+    public Tuple<Collection<ExecutionLineageNode>, 
Collection<ChannelInstance>> evaluate(
+            final ChannelInstance[] inputs,
+            final ChannelInstance[] outputs,
+            final OllamaExecutor ollamaExecutor,
+            final OptimizationContext.OperatorContext operatorContext) {
+        assert inputs.length == this.getNumInputs();
+        assert outputs.length == this.getNumOutputs();
+
+        final Stream<Review> filtered = ((JavaChannelInstance) 
inputs[0]).<Review>provideStream().filter(review -> {
+            try {
+                return OllamaSemanticFilter.isPositiveSentiment2(review);
+            } catch (IOException | InterruptedException e) {
+                throw new RuntimeException("Ollama call failed", e);
+            }
+        });
+        ((StreamChannel.Instance) outputs[0]).accept(filtered);
+
+        return ExecutionOperator.modelLazyExecution(inputs, outputs, 
operatorContext);
+    }
+
+    @Override
+    public String getLoadProfileEstimatorConfigurationKey() {
+        return "wayang.semantic.ollama.model2.load";
+    }
+
+    @Override
+    protected OllamaFilterOperator<Type> newInstance(final DataSetType<Type> 
type, final String prompt) {
+        return new OllamaModel2FilterOperator<>(type, prompt);
+    }
+}
+
+/**
+ * Physical operator for the "model3" Ollama implementation: the only one that 
actually uses the
+ * logical operator's {@link #getPrompt()}.
+ */
+final class OllamaModel3FilterOperator<Type> extends 
OllamaFilterOperator<Type> {
+
+    OllamaModel3FilterOperator(final DataSetType<Type> type, final String 
prompt) {
+        super(type, prompt);
+    }
+
+    @Override
+    @SuppressWarnings("unchecked")
+    public Tuple<Collection<ExecutionLineageNode>, 
Collection<ChannelInstance>> evaluate(
+            final ChannelInstance[] inputs,
+            final ChannelInstance[] outputs,
+            final OllamaExecutor ollamaExecutor,
+            final OptimizationContext.OperatorContext operatorContext) {
+        assert inputs.length == this.getNumInputs();
+        assert outputs.length == this.getNumOutputs();
+
+        final Stream<Review> filtered = ((JavaChannelInstance) 
inputs[0]).<Review>provideStream().filter(review -> {
+            try {
+                return OllamaSemanticFilter.isPositiveSentiment3(review, 
this.getPrompt());
+            } catch (IOException | InterruptedException e) {
+                throw new RuntimeException("Ollama call failed", e);
+            }
+        });
+        ((StreamChannel.Instance) outputs[0]).accept(filtered);
+
+        return ExecutionOperator.modelLazyExecution(inputs, outputs, 
operatorContext);
+    }
+
+    @Override
+    public String getLoadProfileEstimatorConfigurationKey() {
+        return "wayang.semantic.ollama.model3.load";
+    }
+
+    @Override
+    protected OllamaFilterOperator<Type> newInstance(final DataSetType<Type> 
type, final String prompt) {
+        return new OllamaModel3FilterOperator<>(type, prompt);
+    }
+}
+
+abstract class AbstractOllamaFilterMapping implements Mapping {
+
+    private final Class<? extends OllamaFilterOperator> targetOperatorClass;
+
+    AbstractOllamaFilterMapping(final Class<? extends OllamaFilterOperator> 
targetOperatorClass) {
+        this.targetOperatorClass = targetOperatorClass;
+    }
+
+    @Override
+    public Collection<PlanTransformation> getTransformations() {
+        return Collections.singleton(new 
PlanTransformation(this.createSubplanPattern(),
+                this.createReplacementSubplanFactory(), 
OllamaPlatform.getInstance()));
+    }
+
+    private SubplanPattern createSubplanPattern() {
+        return SubplanPattern.createSingleton(new 
OperatorPattern<SemanticFilterOperator<?>>("semantic_filter",
+                new SemanticFilterOperator<>(DataSetType.NONE), false)
+                        .withAdditionalTest(op -> op.getTargetModels() != null)
+                        .withAdditionalTest(op -> 
op.getTargetModels().contains(this.targetOperatorClass)));
+    }
+
+    protected abstract <I> OllamaFilterOperator<I> 
createOperator(DataSetType<I> type, String prompt);
+
+    private <I> ReplacementSubplanFactory createReplacementSubplanFactory() {
+        return new 
ReplacementSubplanFactory.OfSingleOperators<SemanticFilterOperator<I>>((matchedOperator,
 epoch) ->
+                this.<I>createOperator(matchedOperator.getInputType(), 
matchedOperator.getPrompt()).at(epoch));
+    }
+}
+
+final class OllamaModel1FilterMapping extends AbstractOllamaFilterMapping {
+    OllamaModel1FilterMapping() {
+        super(OllamaModel1FilterOperator.class);
+    }
+
+    @Override
+    protected <I> OllamaFilterOperator<I> createOperator(final DataSetType<I> 
type, final String prompt) {
+        return new OllamaModel1FilterOperator<>(type, prompt);
+    }
+}
+
+final class OllamaModel2FilterMapping extends AbstractOllamaFilterMapping {
+    OllamaModel2FilterMapping() {
+        super(OllamaModel2FilterOperator.class);
+    }
+
+    @Override
+    protected <I> OllamaFilterOperator<I> createOperator(final DataSetType<I> 
type, final String prompt) {
+        return new OllamaModel2FilterOperator<>(type, prompt);
+    }
+}
+
+final class OllamaModel3FilterMapping extends AbstractOllamaFilterMapping {
+    OllamaModel3FilterMapping() {
+        super(OllamaModel3FilterOperator.class);
+    }
+
+    @Override
+    protected <I> OllamaFilterOperator<I> createOperator(final DataSetType<I> 
type, final String prompt) {
+        return new OllamaModel3FilterOperator<>(type, prompt);
+    }
+}
+
+final class OllamaSemanticFilter {
+    private static volatile String OLLAMA_API_URL = 
"http://ollama:11434/api/generate";;
+    private static final String MODEL_NAME = "tinyllama";
+    private static final HttpClient httpClient = HttpClient.newHttpClient();
+
+    /**
+     * Points calls at a different Ollama-compatible endpoint, e.g. an 
in-process mock server for
+     * tests, so the real HTTP wiring gets exercised without depending on an 
external Ollama instance.
+     */
+    static void setApiUrl(final String apiUrl) {
+        OLLAMA_API_URL = apiUrl;
+    }
+
+    public static boolean isPositiveSentiment(final Review review) throws 
IOException, InterruptedException {
+        final String prompt = String.format("Analyze the sentiment of this 
movie review. "
+                + "Reply with only 'POSITIVE' or 'NEGATIVE'.\n\n" + "Review: 
%s", review.getReviewText());
+        final String response = callOllama(prompt);
+        return response.contains("POSITIVE");
+    }
+
+    public static boolean isPositiveSentiment2(final Review review) throws 
IOException, InterruptedException {
+        final String prompt = String.format(
+                "Analyze the sentiment of this movie review, words like love, 
fantastic and great are positive modifiers. "
+                        + "Reply with only 'POSITIVE' or 'NEGATIVE'.\n\n" + 
"Review: %s",
+                review.getReviewText());
+        final String response = callOllama(prompt);
+        return response.contains("POSITIVE");
+    }
+
+    public static boolean isPositiveSentiment3(final Review review, final 
String prompt)
+            throws IOException, InterruptedException {
+        final String response = callOllama(prompt + " | " + 
review.getReviewText());
+        return response.contains("POSITIVE");
+    }
+
+    private static String callOllama(final String prompt) throws IOException, 
InterruptedException {
+        final String requestBody = String.format("{\"model\": \"%s\", 
\"prompt\": \"%s\", \"stream\": false}",
+                MODEL_NAME, escapeJson(prompt));
+
+        final HttpRequest request = 
HttpRequest.newBuilder().uri(URI.create(OLLAMA_API_URL))
+                .header("Content-Type", 
"application/json").POST(HttpRequest.BodyPublishers.ofString(requestBody))
+                .timeout(Duration.ofMinutes(2)).build();
+
+        final HttpResponse<String> response = httpClient.send(request, 
HttpResponse.BodyHandlers.ofString());
+
+        if (response.statusCode() != 200) {
+            throw new IOException("Ollama API error: " + response.body());
+        }
+
+        return parseOllamaResponse(response.body());
+    }
+
+    private static String parseOllamaResponse(final String jsonResponse) {
+        int startIdx = jsonResponse.indexOf("\"response\":\"");
+
+        if (startIdx == -1)
+            return "";
+
+        startIdx += "\"response\":\"".length();
+
+        final int endIdx = jsonResponse.indexOf("\"", startIdx);
+
+        return jsonResponse.substring(startIdx, endIdx);
+    }
+
+    private static String escapeJson(final String str) {
+        return str.replace("\"", "\\\"").replace("\n", "\\n").replace("\r", 
"\\r");
+    }
+}

Reply via email to