diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java index eadfd7458c140..a4bc06ccad187 100644 --- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java +++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/AnalyticsPlugin.java @@ -30,6 +30,7 @@ import org.opensearch.cluster.service.ClusterService; import org.opensearch.common.inject.Module; import org.opensearch.common.inject.TypeLiteral; +import org.opensearch.common.settings.Setting; import org.opensearch.core.action.ActionResponse; import org.opensearch.core.common.io.stream.NamedWriteableRegistry; import org.opensearch.core.xcontent.NamedXContentRegistry; @@ -62,6 +63,14 @@ public class AnalyticsPlugin extends Plugin implements ExtensiblePlugin, ActionP private static final Logger logger = LogManager.getLogger(AnalyticsPlugin.class); + public static final Setting COORDINATOR_BUFFER_LIMIT = Setting.longSetting( + "analytics.coordinator.buffer_limit", + 256L * 1024 * 1024, + 0L, + Setting.Property.NodeScope, + Setting.Property.Dynamic + ); + /** * Creates a new analytics engine hub plugin. */ @@ -129,6 +138,11 @@ public Collection createGuiceModules() { return List.of(new ActionHandler<>(AnalyticsQueryAction.INSTANCE, DefaultPlanExecutor.class)); } + @Override + public List> getSettings() { + return List.of(COORDINATOR_BUFFER_LIMIT); + } + @Override public void close() { if (searchService != null) { diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java index 09bd28b0294ea..d4df714ae863f 100644 --- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java +++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/DefaultPlanExecutor.java @@ -8,6 +8,7 @@ package org.opensearch.analytics.exec; +import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.metadata.JaninoRelMetadataProvider; @@ -18,6 +19,7 @@ import org.opensearch.action.support.ActionFilters; import org.opensearch.action.support.HandledTransportAction; import org.opensearch.action.support.TimeoutTaskCancellationUtility; +import org.opensearch.analytics.AnalyticsPlugin; import org.opensearch.analytics.EngineContext; import org.opensearch.analytics.exec.action.AnalyticsQueryAction; import org.opensearch.analytics.exec.task.AnalyticsQueryTask; @@ -73,7 +75,10 @@ public class DefaultPlanExecutor extends HandledTransportAction perQueryBufferLimit = v); } @Override @@ -144,7 +152,23 @@ private void executeInternal(RelNode logicalFragment, ActionListener> batchesListener = ActionListener.runAfter( ActionListener.wrap(batches -> listener.onResponse(batchesToRows(batches)), listener::onFailure), diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryContext.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryContext.java index d3bb142004eb9..0d708d19bbc5d 100644 --- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryContext.java +++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryContext.java @@ -13,7 +13,6 @@ import org.opensearch.analytics.backend.AnalyticsOperationListener; import org.opensearch.analytics.exec.task.AnalyticsQueryTask; import org.opensearch.analytics.planner.dag.QueryDAG; -import org.opensearch.arrow.memory.ArrowAllocatorService; import java.util.List; import java.util.concurrent.Executor; @@ -31,30 +30,24 @@ public class QueryContext { // TODO: make configurable via cluster setting (like search.max_concurrent_shard_requests) private static final int DEFAULT_MAX_CONCURRENT_SHARD_REQUESTS = 5; - /** Default per-query memory limit for Arrow allocations (256 MB). */ - private static final long DEFAULT_PER_QUERY_MEMORY_LIMIT = 256L * 1024 * 1024; - private final QueryDAG dag; private final Executor searchExecutor; private final AnalyticsQueryTask parentTask; private final int maxConcurrentShardRequests; - private final long perQueryMemoryLimit; private final List operationListeners; - private final ArrowAllocatorService allocatorService; - private volatile BufferAllocator bufferAllocator; + private final BufferAllocator allocator; + private final boolean ownsAllocator; private volatile ExecutorService localTaskExecutor; private boolean closed; // guarded by `this` - public QueryContext(QueryDAG dag, Executor searchExecutor, AnalyticsQueryTask parentTask, ArrowAllocatorService allocatorService) { - this( - dag, - searchExecutor, - parentTask, - DEFAULT_MAX_CONCURRENT_SHARD_REQUESTS, - DEFAULT_PER_QUERY_MEMORY_LIMIT, - List.of(), - allocatorService - ); + public QueryContext( + QueryDAG dag, + Executor searchExecutor, + AnalyticsQueryTask parentTask, + BufferAllocator allocator, + boolean ownsAllocator + ) { + this(dag, searchExecutor, parentTask, DEFAULT_MAX_CONCURRENT_SHARD_REQUESTS, List.of(), allocator, ownsAllocator); } /** Full-parameter constructor. Private; tests use {@link #forTest} factories. */ @@ -63,17 +56,17 @@ private QueryContext( Executor searchExecutor, AnalyticsQueryTask parentTask, int maxConcurrentShardRequests, - long perQueryMemoryLimit, List operationListeners, - ArrowAllocatorService allocatorService + BufferAllocator allocator, + boolean ownsAllocator ) { this.dag = dag; this.searchExecutor = searchExecutor; this.parentTask = parentTask; this.maxConcurrentShardRequests = maxConcurrentShardRequests; - this.perQueryMemoryLimit = perQueryMemoryLimit; this.operationListeners = operationListeners; - this.allocatorService = allocatorService; + this.allocator = allocator; + this.ownsAllocator = ownsAllocator; } public QueryDAG dag() { @@ -101,22 +94,8 @@ public List operationListeners() { return operationListeners; } - /** Lazy per-query allocator (child of shared root) with {@link #perQueryMemoryLimit}. */ public BufferAllocator bufferAllocator() { - BufferAllocator alloc = bufferAllocator; - if (alloc == null) { - synchronized (this) { - alloc = bufferAllocator; - if (alloc == null) { - if (closed) { - throw new IllegalStateException("QueryContext closed for query " + dag.queryId()); - } - alloc = allocatorService.newChildAllocator("query-" + dag.queryId(), perQueryMemoryLimit); - bufferAllocator = alloc; - } - } - } - return alloc; + return allocator; } /** Lazy per-query virtual-thread executor for LOCAL tasks. */ @@ -139,14 +118,17 @@ public ExecutorService localTaskExecutor() { return exec; } + boolean ownsAllocator() { + return ownsAllocator; + } + /** Idempotent. Serialised with lazy-init accessors; post-close accessors throw. */ public void close() { synchronized (this) { if (closed) return; closed = true; - if (bufferAllocator != null) { - bufferAllocator.close(); - bufferAllocator = null; + if (ownsAllocator) { + allocator.close(); } if (localTaskExecutor != null) { localTaskExecutor.shutdown(); @@ -157,27 +139,7 @@ public void close() { // ─── Test factories ──────────────────────────────────────────────── - /** Test-only: wraps a fresh {@link RootAllocator} as an {@link ArrowAllocatorService}. */ - private static ArrowAllocatorService testAllocatorService() { - return new ArrowAllocatorService() { - private final RootAllocator root = new RootAllocator(Long.MAX_VALUE); - - @Override - public BufferAllocator newChildAllocator(String name, long limit) { - return root.newChildAllocator(name, 0, limit); - } - - @Override - public long getAllocatedMemory() { - return root.getAllocatedMemory(); - } - - @Override - public long getPeakMemoryAllocation() { - return root.getPeakMemoryAllocation(); - } - }; - } + private static final RootAllocator TEST_ROOT = new RootAllocator(Long.MAX_VALUE); /** Creates a test context with a synchronous executor. */ public static QueryContext forTest(QueryDAG dag, AnalyticsQueryTask parentTask) { @@ -186,14 +148,15 @@ public static QueryContext forTest(QueryDAG dag, AnalyticsQueryTask parentTask) /** Creates a test context with a synchronous executor and the supplied operation listeners. */ public static QueryContext forTest(QueryDAG dag, AnalyticsQueryTask parentTask, List operationListeners) { + BufferAllocator testAllocator = TEST_ROOT.newChildAllocator("test-" + dag.queryId(), 0, Long.MAX_VALUE); return new QueryContext( dag, Runnable::run, parentTask, DEFAULT_MAX_CONCURRENT_SHARD_REQUESTS, - Long.MAX_VALUE, operationListeners, - testAllocatorService() + testAllocator, + true ); } } diff --git a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryExecution.java b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryExecution.java index 5458ba7518b66..07a9d728aba97 100644 --- a/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryExecution.java +++ b/sandbox/plugins/analytics-engine/src/main/java/org/opensearch/analytics/exec/QueryExecution.java @@ -8,6 +8,7 @@ package org.opensearch.analytics.exec; +import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -118,9 +119,21 @@ public void close() { if (closed.compareAndSet(false, true) == false) return; runQuietly("terminal sink close", this::closeTerminalSink); // TODO: Re-evaluate this per query child allocator + logAllocatorState(); runQuietly("query context close", config::close); } + private void logAllocatorState() { + if (!config.ownsAllocator()) return; + BufferAllocator allocator = config.bufferAllocator(); + long allocated = allocator.getAllocatedMemory(); + if (allocated > 0) { + logger.warn("[query-{}] Arrow allocator closing with {}B still allocated — potential leak", config.queryId(), allocated); + } else { + logger.debug("[query-{}] Arrow allocator closed cleanly", config.queryId()); + } + } + // ─── Internal: query-level state machine ───────────────────────────── /** On terminal transition: fires user listener exactly once + runs {@link #close()}. */