This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new a333bdc71e [spark] Support query operator in full_text_search TVF
(#8281)
a333bdc71e is described below
commit a333bdc71edad4ee609251926344b6cc4b385612
Author: QuakeWang <[email protected]>
AuthorDate: Fri Jun 19 08:35:48 2026 +0800
[spark] Support query operator in full_text_search TVF (#8281)
Spark `full_text_search` TVF only accepted the original four arguments
and always built the default OR query, while the underlying full-text
predicate and hybrid full-text route already supported `query_operator`.
This PR adds an optional `query_operator` argument to Spark
`full_text_search`, keeps `or` as the default for backward
compatibility, and reuses `FullTextSearch` for `or` / `and` validation
and normalization.
---
.../multimodal-table/global-index/full-text.mdx | 5 +++-
.../TestFullTextGlobalIndexReader.java | 8 ++++--
.../plans/logical/PaimonTableValuedFunctions.scala | 18 ++++++++----
.../plans/logical/VectorSearchQueryTest.scala | 32 ++++++++++++++++++++++
.../paimon/spark/sql/FullTextSearchTest.scala | 21 ++++++++------
5 files changed, 68 insertions(+), 16 deletions(-)
diff --git a/docs/docs/multimodal-table/global-index/full-text.mdx
b/docs/docs/multimodal-table/global-index/full-text.mdx
index 1e65d490a6..7c3577fa87 100644
--- a/docs/docs/multimodal-table/global-index/full-text.mdx
+++ b/docs/docs/multimodal-table/global-index/full-text.mdx
@@ -125,8 +125,11 @@ The query operator controls how multiple query terms are
matched:
<TabItem value="spark-sql" label="Spark SQL">
```sql
--- Search for top-10 documents matching the query
+-- Search for top-10 documents matching any query term. The default query
operator is 'or'.
SELECT * FROM full_text_search('my_table', 'content', 'paimon lake format',
10);
+
+-- Search for top-10 documents matching all query terms.
+SELECT * FROM full_text_search('my_table', 'content', 'paimon lake format',
10, 'and');
```
</TabItem>
diff --git
a/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexReader.java
b/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexReader.java
index 7e9684f267..d4ff47cca4 100644
---
a/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexReader.java
+++
b/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexReader.java
@@ -80,13 +80,14 @@ public class TestFullTextGlobalIndexReader implements
GlobalIndexReader {
}
String[] queryTerms = queryText.toLowerCase(Locale.ROOT).split("\\s+");
+ boolean requireAllTerms = "and".equals(fullTextSearch.queryOperator());
// Min-heap: smallest score at head, so we evict the weakest candidate.
PriorityQueue<ScoredRow> topK =
new PriorityQueue<>(effectiveK + 1,
Comparator.comparingDouble(s -> s.score));
for (int i = 0; i < count; i++) {
- float score = computeScore(documents[i], queryTerms);
+ float score = computeScore(documents[i], queryTerms,
requireAllTerms);
if (score <= 0) {
continue;
}
@@ -113,12 +114,15 @@ public class TestFullTextGlobalIndexReader implements
GlobalIndexReader {
Optional.of(ScoredGlobalIndexResult.create(resultBitmap,
scoreMap::get)));
}
- private static float computeScore(String document, String[] queryTerms) {
+ private static float computeScore(
+ String document, String[] queryTerms, boolean requireAllTerms) {
String lowerDoc = document.toLowerCase(Locale.ROOT);
float score = 0;
for (String term : queryTerms) {
if (lowerDoc.contains(term)) {
score += 1.0f / queryTerms.length;
+ } else if (requireAllTerms) {
+ return 0;
}
}
return score;
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonTableValuedFunctions.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonTableValuedFunctions.scala
index 46af6bb38a..05f70bcb37 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonTableValuedFunctions.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/catalyst/plans/logical/PaimonTableValuedFunctions.scala
@@ -648,13 +648,14 @@ case class HybridSearchQuery(override val args:
Seq[Expression])
/**
* Plan for the [[FULL_TEXT_SEARCH]] table-valued function.
*
- * Usage: full_text_search(table_name, column_name, query_text, limit)
+ * Usage: full_text_search(table_name, column_name, query_text, limit[,
query_operator])
* - table_name: the Paimon table to search
* - column_name: the text column name
* - query_text: the query text string
* - limit: the number of top results to return
+ * - query_operator: optional query operator, supported values are 'or' and
'and'
*
- * Example: SELECT * FROM full_text_search('T', 'content', 'hello world', 10)
+ * Example: SELECT * FROM full_text_search('T', 'content', 'hello world', 10,
'and')
*/
case class FullTextSearchQuery(override val args: Seq[Expression])
extends PaimonTableValueFunction(FULL_TEXT_SEARCH) {
@@ -667,9 +668,10 @@ case class FullTextSearchQuery(override val args:
Seq[Expression])
def createFullTextSearch(
innerTable: InnerTable,
argsWithoutTable: Seq[Expression]): FullTextSearch = {
- if (argsWithoutTable.size != 3) {
+ if (argsWithoutTable.size != 3 && argsWithoutTable.size != 4) {
throw new RuntimeException(
- s"$FULL_TEXT_SEARCH needs three parameters after table_name:
column_name, query_text, limit. " +
+ s"$FULL_TEXT_SEARCH needs three or four parameters after table_name: "
+
+ s"column_name, query_text, limit[, query_operator]. " +
s"Got ${argsWithoutTable.size} parameters after table_name."
)
}
@@ -681,6 +683,12 @@ case class FullTextSearchQuery(override val args:
Seq[Expression])
}
val queryText = argsWithoutTable(1).eval().toString
val limit = parsePositiveLimit(argsWithoutTable(2).eval())
- new FullTextSearch(queryText, limit, columnName)
+ val queryOperator =
+ if (argsWithoutTable.size == 4) {
+ VectorSearchQuery(Seq.empty).extractString(argsWithoutTable(3))
+ } else {
+ "or"
+ }
+ new FullTextSearch(queryText, limit, columnName, queryOperator)
}
}
diff --git
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/catalyst/plans/logical/VectorSearchQueryTest.scala
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/catalyst/plans/logical/VectorSearchQueryTest.scala
index b3b5811f2b..23427a54e6 100644
---
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/catalyst/plans/logical/VectorSearchQueryTest.scala
+++
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/catalyst/plans/logical/VectorSearchQueryTest.scala
@@ -137,6 +137,38 @@ class VectorSearchQueryTest extends AnyFunSuite {
assert(search.routes().get(0).weight() == 1.5f)
}
+ test("create full-text search with default query operator") {
+ val search = FullTextSearchQuery(Seq.empty).createFullTextSearch(
+ innerTable,
+ Seq(Literal("content"), Literal("paimon lake"), Literal(10)))
+
+ assert(search.fieldName() == "content")
+ assert(search.queryText() == "paimon lake")
+ assert(search.limit() == 10)
+ assert(search.queryOperator() == "or")
+ }
+
+ test("create full-text search with explicit query operator") {
+ val search = FullTextSearchQuery(Seq.empty).createFullTextSearch(
+ innerTable,
+ Seq(Literal("content"), Literal("paimon lake"), Literal(10),
Literal("and")))
+
+ assert(search.fieldName() == "content")
+ assert(search.queryText() == "paimon lake")
+ assert(search.limit() == 10)
+ assert(search.queryOperator() == "and")
+ }
+
+ test("reject invalid full-text search query operator") {
+ val exception = intercept[IllegalArgumentException] {
+ FullTextSearchQuery(Seq.empty).createFullTextSearch(
+ innerTable,
+ Seq(Literal("content"), Literal("paimon lake"), Literal(10),
Literal("xor")))
+ }
+
+ assert(exception.getMessage.contains("Query operator must be 'or' or
'and'"))
+ }
+
test("reject hybrid search query map") {
val exception = intercept[RuntimeException] {
HybridSearchQuery(Seq.empty).createHybridSearch(
diff --git
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FullTextSearchTest.scala
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FullTextSearchTest.scala
index 8f6d06fd75..179cb32645 100644
---
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FullTextSearchTest.scala
+++
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/sql/FullTextSearchTest.scala
@@ -137,7 +137,7 @@ class FullTextSearchTest extends PaimonSparkTestBase {
}
}
- test("full-text search - multi-term query") {
+ test("full-text search - multi-term query operators") {
withTable("T") {
spark.sql("""
|CREATE TABLE T (id INT, content STRING)
@@ -162,18 +162,23 @@ class FullTextSearchTest extends PaimonSparkTestBase {
s"CALL sys.create_global_index(table => 'test.T', index_column =>
'content', index_type => '$indexType')")
.collect()
- // Query "Paimon search" - rows 1 and 2 match both terms
- val result = spark
+ val defaultOrResult = spark
+ .sql("""
+ |SELECT id FROM full_text_search('T', 'content', 'Paimon
search', 5)
+ |ORDER BY id
+ |""".stripMargin)
+ .collect()
+
+ assert(defaultOrResult.map(_.getInt(0)).toSeq == Seq(0, 1, 2, 3))
+
+ val explicitAndResult = spark
.sql("""
- |SELECT id FROM full_text_search('T', 'content', 'Paimon
search', 2)
+ |SELECT id FROM full_text_search('T', 'content', 'Paimon
search', 5, 'and')
|ORDER BY id
|""".stripMargin)
.collect()
- assert(result.length == 2)
- val ids = result.map(_.getInt(0)).toSet
- assert(ids.contains(1))
- assert(ids.contains(2))
+ assert(explicitAndResult.map(_.getInt(0)).toSeq == Seq(1, 2))
}
}