Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,6 @@ public enum Key {

/** Query Settings. */
FIELD_TYPE_TOLERANCE("plugins.query.field_type_tolerance"),
PARTIAL_RESULT_ON_MAPPING_CONFLICT("plugins.query.partial_result.on_mapping_conflict.enabled"),

/** Common Settings for SQL and PPL. */
QUERY_MEMORY_LIMIT("plugins.query.memory_limit"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,32 +65,13 @@ public class CalcitePlanContext {
public static final ThreadLocal<String> executionPool = new ThreadLocal<>();

/**
* Non-fatal warnings raised during planning (e.g. a partial result over a subset of indices) to
* be attached to the query response by the execution engine. Drained in {@code
* OpenSearchExecutionEngine.buildResultSet} and cleared with the other lifecycle signals so it
* never leaks onto the next query on a pooled worker thread.
* Non-fatal warnings raised while planning or executing a query, to be attached to its response.
* Drained in {@code OpenSearchExecutionEngine.buildResultSet} and cleared with the other
* lifecycle signals so it never leaks onto the next query on a pooled worker thread.
*/
private static final ThreadLocal<List<Warning>> pendingWarnings =
ThreadLocal.withInitial(ArrayList::new);

/**
* Whether the current query's response format can carry a warnings channel. Set on the worker
* thread from the plan (see {@code QueryPlan#execute}) rather than the transport thread, so the
* partial-result gate survives the transport→worker handoff — the security plugin's interceptor
* drops Log4j {@code ThreadContext}, which is where this used to live. Cleared per query.
*/
private static final ThreadLocal<Boolean> warningsSupported =
ThreadLocal.withInitial(() -> false);

/**
* Per-request partial-result override, carried off Log4j {@code ThreadContext} onto the plan (see
* {@code QueryPlan#execute}) for the same reason as {@link #warningsSupported}: the security
* plugin's interceptor drops {@code ThreadContext} on the transport→worker handoff. {@code null}
* defers to the cluster setting; {@code true}/{@code false} force partial mode on/off for this
* query. Cleared per query.
*/
private static final ThreadLocal<Boolean> partialResultOverride = new ThreadLocal<>();

/** Thread-local switch that tells whether the current query prefers legacy behavior. */
private static final ThreadLocal<Boolean> legacyPreferredFlag =
ThreadLocal.withInitial(() -> true);
Expand Down Expand Up @@ -279,44 +260,13 @@ public static void clearTimewrapSignals() {
timewrapSeries.set(null);
executionPool.set(null);
pendingWarnings.remove();
warningsSupported.set(false);
partialResultOverride.remove();
}

/** Records a non-fatal warning to be attached to the response for the current query. */
public static void addWarning(Warning warning) {
pendingWarnings.get().add(warning);
}

/** Records whether the current query's response format can surface warnings. */
public static void setWarningsSupported(boolean supported) {
warningsSupported.set(supported);
}

/**
* @return whether the current query's response format can surface warnings; false when unset, so
* a caller that never declared support cannot get a silent partial result.
*/
public static boolean isWarningsSupported() {
return warningsSupported.get();
}

/**
* Records the per-request partial-result override for the current query. {@code null} defers to
* the cluster setting; {@code true}/{@code false} force partial mode on/off.
*/
public static void setPartialResultOverride(Boolean override) {
partialResultOverride.set(override);
}

/**
* @return the per-request partial-result override, or {@code null} to defer to the cluster
* setting.
*/
public static Boolean getPartialResultOverride() {
return partialResultOverride.get();
}

/**
* Returns and clears the warnings collected for the current query, de-duplicated by value. The
* planner may fire a rule that raises a warning more than once for equivalent plan alternatives,
Expand All @@ -342,24 +292,18 @@ public static class ThreadLocalSnapshot {
final String timewrapUnitName;
final String timewrapSeries;
final String executionPool;
final boolean warningsSupported;
final Boolean partialResultOverride;

private ThreadLocalSnapshot(
boolean skipEncoding,
boolean stripNullColumns,
String timewrapUnitName,
String timewrapSeries,
String executionPool,
boolean warningsSupported,
Boolean partialResultOverride) {
String executionPool) {
this.skipEncoding = skipEncoding;
this.stripNullColumns = stripNullColumns;
this.timewrapUnitName = timewrapUnitName;
this.timewrapSeries = timewrapSeries;
this.executionPool = executionPool;
this.warningsSupported = warningsSupported;
this.partialResultOverride = partialResultOverride;
}
}

Expand All @@ -370,9 +314,7 @@ public static ThreadLocalSnapshot snapshotThreadLocals() {
stripNullColumns.get(),
timewrapUnitName.get(),
timewrapSeries.get(),
executionPool.get(),
warningsSupported.get(),
partialResultOverride.get());
executionPool.get());
}

/** Restore thread-local state from a snapshot. */
Expand All @@ -382,8 +324,6 @@ public static void restoreThreadLocals(ThreadLocalSnapshot snapshot) {
timewrapUnitName.set(snapshot.timewrapUnitName);
timewrapSeries.set(snapshot.timewrapSeries);
executionPool.set(snapshot.executionPool);
warningsSupported.set(snapshot.warningsSupported);
partialResultOverride.set(snapshot.partialResultOverride);
}

public void pushForeachBindings(
Expand Down
11 changes: 5 additions & 6 deletions core/src/main/java/org/opensearch/sql/executor/Warning.java
Original file line number Diff line number Diff line change
Expand Up @@ -9,17 +9,16 @@

/**
* A non-fatal notice attached to an otherwise-successful query response. Carried through the
* response path so consumers can distinguish a correct-but-noteworthy result (e.g. a partial result
* over a subset of indices) from a plain success, without turning it into an error.
* response path so consumers can distinguish a correct-but-noteworthy result from a plain success,
* without turning it into an error.
*/
@Data
public class Warning {

/**
* The result is complete for the indices it covers but omits one or more indices that could not
* be served (e.g. a mapping conflict that prevents aggregation pushdown). This is a cross-surface
* contract: consumers such as OpenSearch Dashboards branch on this {@code type} value, so it must
* not change without coordinating those consumers.
* The response does not cover everything the query asked for. This is a cross-surface contract:
* consumers such as OpenSearch Dashboards branch on this {@code type} value, so it must not
* change without coordinating those consumers.
*/
public static final String TYPE_PARTIAL_RESULT = "PARTIAL_RESULT";

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@

import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.Setter;
import org.opensearch.sql.ast.statement.ExplainMode;
import org.opensearch.sql.common.response.ResponseListener;
import org.opensearch.sql.executor.ExecutionEngine;
Expand All @@ -24,22 +23,6 @@ public abstract class AbstractPlan {

@Getter protected final QueryType queryType;

/**
* Whether the response format can carry a warnings channel. Set from the request on the transport
* thread and read on the worker (see {@code QueryPlan#execute}), so a feature that returns a
* partial result never silently drops data into a warnings-incapable format across a thread
* handoff. Defaults to false.
*/
@Getter @Setter private boolean warningsSupported = false;

/**
* Per-request partial-result override, carried from the request the same way as {@link
* #warningsSupported}. {@code null} defers to the cluster setting; {@code true}/{@code false}
* force partial mode on/off. Set on the transport thread, applied on the worker (see {@code
* QueryPlan#execute}) so it survives the security transport→worker handoff.
*/
@Getter @Setter private Boolean partialResultOverride = null;

/** Start query execution. */
public abstract void execute();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
import org.opensearch.sql.ast.tree.HighlightConfig;
import org.opensearch.sql.ast.tree.Paginate;
import org.opensearch.sql.ast.tree.UnresolvedPlan;
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.common.response.ResponseListener;
import org.opensearch.sql.executor.ExecutionEngine;
import org.opensearch.sql.executor.QueryId;
Expand Down Expand Up @@ -106,11 +105,6 @@ public QueryPlan(

@Override
public void execute() {
// Runs on the worker thread; carry warnings support and the per-request partial-result override
// from the request off the plan so the partial-result gate reads them without depending on
// Log4j ThreadContext (dropped under security on the transport→worker handoff).
CalcitePlanContext.setWarningsSupported(isWarningsSupported());
CalcitePlanContext.setPartialResultOverride(getPartialResultOverride());
if (pageSize.isPresent()) {
queryService.execute(
new Paginate(pageSize.get(), plan),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@

package org.opensearch.sql.executor.execution;

import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
Expand All @@ -17,7 +16,6 @@
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;

import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.commons.lang3.NotImplementedException;
import org.junit.jupiter.api.DisplayNameGeneration;
import org.junit.jupiter.api.DisplayNameGenerator;
Expand All @@ -27,7 +25,6 @@
import org.mockito.junit.jupiter.MockitoExtension;
import org.opensearch.sql.ast.statement.ExplainMode;
import org.opensearch.sql.ast.tree.UnresolvedPlan;
import org.opensearch.sql.calcite.CalcitePlanContext;
import org.opensearch.sql.common.response.ResponseListener;
import org.opensearch.sql.executor.DefaultExecutionEngine;
import org.opensearch.sql.executor.ExecutionEngine;
Expand Down Expand Up @@ -61,53 +58,6 @@ public void execute_no_page_size() {
verify(queryService, times(1)).execute(any(), any(), any(), anyBoolean(), any());
}

@Test
public void warnings_supported_flag_reaches_execution_thread() throws InterruptedException {
QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener);
query.setWarningsSupported(true);

// Configure the plan here but run execute() on another thread, mirroring the transport->worker
// handoff. The flag must ride the plan object, not a thread-local the handoff can drop.
AtomicBoolean defaultedBeforeExecute = new AtomicBoolean(true);
AtomicBoolean seenOnWorker = new AtomicBoolean(false);
Thread worker =
new Thread(
() -> {
defaultedBeforeExecute.set(CalcitePlanContext.isWarningsSupported());
query.execute();
seenOnWorker.set(CalcitePlanContext.isWarningsSupported());
});
worker.start();
worker.join();

assertFalse(
defaultedBeforeExecute.get(), "worker thread should default to no warnings support");
assertTrue(seenOnWorker.get(), "execute() must carry warningsSupported onto the worker thread");
verify(queryService, times(1)).execute(any(), any(), any(), anyBoolean(), any());
}

@Test
public void warnings_unsupported_plan_resets_flag_on_reused_worker_thread()
throws InterruptedException {
// Plan defaults to warningsSupported=false.
QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener);

AtomicBoolean seenOnWorker = new AtomicBoolean(true);
Thread worker =
new Thread(
() -> {
// Simulate a pooled worker left "supported" by a prior query.
CalcitePlanContext.setWarningsSupported(true);
query.execute();
seenOnWorker.set(CalcitePlanContext.isWarningsSupported());
});
worker.start();
worker.join();

assertFalse(
seenOnWorker.get(), "a warnings-unsupported plan must reset the flag on a reused thread");
}

@Test
public void explain_no_page_size() {
QueryPlan query = new QueryPlan(queryId, queryType, plan, queryService, queryListener);
Expand Down
47 changes: 0 additions & 47 deletions docs/user/admin/settings.rst
Original file line number Diff line number Diff line change
Expand Up @@ -358,53 +358,6 @@ Result set::
}
}

plugins.query.partial_result.on_mapping_conflict.enabled [Experimental]
=======================================================================

Version
-------
Since 3.9

Description
-----------

This setting is experimental; its name, values, and default may change in a future release. Controls how an aggregation behaves when its group-by field is mapped inconsistently across the queried indices -- for example ``keyword`` in some indices of a wildcard pattern and ``text`` (without a ``.keyword`` sub-field) in others. Such a field collapses to ``text``-without-``.keyword`` across the pattern, which has no doc values, so the aggregation cannot be pushed down natively and instead runs as a per-document script over ``_source`` -- correct, but a full scan of every document.

When this setting is ``false`` (the default), that complete-but-slow result is returned. When set to ``true``, the aggregation is pushed down over only the subset of indices where the field is aggregatable, and the response carries a ``PARTIAL_RESULT`` warning naming the excluded indices and the remedy (map the field as ``keyword`` everywhere). The result is therefore **partial** -- documents in the excluded indices are not counted -- so the setting is off by default and only takes effect for response formats that can surface the warning (the JSON format; CSV/raw/visualization responses fall through to the complete result rather than silently dropping data).

The behavior can also be overridden per request with the ``partial_result`` boolean field in the query body, which takes precedence over this cluster setting. Here is an example enabling it at the cluster level::

>> curl -H 'Content-Type: application/json' -X PUT localhost:9200/_plugins/_query/settings -d '{
"transient" : {
"plugins.query.partial_result.on_mapping_conflict.enabled" : true
}
}'

Result set::

{
"acknowledged" : true,
"persistent" : { },
"transient" : {
"plugins" : {
"query" : {
"partial_result" : {
"on_mapping_conflict" : {
"enabled" : "true"
}
}
}
}
}
}

Per-request override example, opting a single query into a partial result regardless of the cluster setting::

>> curl -H 'Content-Type: application/json' -X POST localhost:9200/_plugins/_ppl -d '{
"query" : "source=logs-* | stats count() by service",
"partial_result" : true
}'

plugins.query.buckets
=====================

Expand Down
Loading
Loading