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
@@ -0,0 +1,28 @@
/*
* Licensed 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 com.facebook.presto;

import com.facebook.presto.spi.PrestoException;
import io.airlift.units.DataSize;

import static com.facebook.presto.spi.StandardErrorCode.EXCEEDED_SCAN_RAW_BYTES_READ_LIMIT;

public class ExceededScanLimitException
extends PrestoException
{
public ExceededScanLimitException(DataSize limit)
{
super(EXCEEDED_SCAN_RAW_BYTES_READ_LIMIT, "Query has exceeded Scan Raw Bytes Read Limit of " + limit.toString());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
import static com.facebook.presto.common.type.VarcharType.VARCHAR;
import static com.facebook.presto.spi.StandardErrorCode.INVALID_SESSION_PROPERTY;
import static com.facebook.presto.spi.session.PropertyMetadata.booleanProperty;
import static com.facebook.presto.spi.session.PropertyMetadata.dataSizeProperty;
import static com.facebook.presto.spi.session.PropertyMetadata.doubleProperty;
import static com.facebook.presto.spi.session.PropertyMetadata.integerProperty;
import static com.facebook.presto.spi.session.PropertyMetadata.stringProperty;
Expand Down Expand Up @@ -90,6 +91,7 @@ public final class SystemSessionProperties
public static final String QUERY_MAX_RUN_TIME = "query_max_run_time";
public static final String RESOURCE_OVERCOMMIT = "resource_overcommit";
public static final String QUERY_MAX_CPU_TIME = "query_max_cpu_time";
public static final String QUERY_MAX_SCAN_RAW_INPUT_BYTES = "query_max_scan_raw_input_bytes";
public static final String QUERY_MAX_STAGE_COUNT = "query_max_stage_count";
public static final String REDISTRIBUTE_WRITES = "redistribute_writes";
public static final String SCALE_WRITERS = "scale_writers";
Expand Down Expand Up @@ -424,6 +426,11 @@ public SystemSessionProperties(
"Use resources which are not guaranteed to be available to the query",
false,
false),
dataSizeProperty(
QUERY_MAX_SCAN_RAW_INPUT_BYTES,
"Maximum scan raw input bytes of a query",
queryManagerConfig.getQueryMaxScanRawInputBytes(),
false),
integerProperty(
QUERY_MAX_STAGE_COUNT,
"Temporary: Maximum number of stages a query can have",
Expand Down Expand Up @@ -1115,6 +1122,11 @@ public static Duration getQueryMaxCpuTime(Session session)
return session.getSystemProperty(QUERY_MAX_CPU_TIME, Duration.class);
}

public static DataSize getQueryMaxScanRawInputBytes(Session session)
{
return session.getSystemProperty(QUERY_MAX_SCAN_RAW_INPUT_BYTES, DataSize.class);
}

public static boolean isSpillEnabled(Session session)
{
return session.getSystemProperty(SPILL_ENABLED, Boolean.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,12 @@ public Duration getTotalCpuTime()
return new Duration(0, NANOSECONDS);
}

@Override
public DataSize getRawInputDataSize()
{
return DataSize.succinctBytes(0);
}

@Override
public BasicQueryInfo getBasicQueryInfo()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@ public interface QueryExecution

Duration getTotalCpuTime();

DataSize getRawInputDataSize();

DataSize getUserMemoryReservation();

DataSize getTotalMemoryReservation();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import com.facebook.airlift.configuration.DefunctConfig;
import com.facebook.airlift.configuration.LegacyConfig;
import com.facebook.presto.connector.system.GlobalSystemConnector;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;
import io.airlift.units.MinDuration;

Expand All @@ -27,6 +28,8 @@

import java.util.concurrent.TimeUnit;

import static io.airlift.units.DataSize.Unit.PETABYTE;

@DefunctConfig({
"query.max-pending-splits-per-node",
"query.queue-config-file",
Expand Down Expand Up @@ -66,6 +69,8 @@ public class QueryManagerConfig
private Duration queryMaxExecutionTime = new Duration(100, TimeUnit.DAYS);
private Duration queryMaxCpuTime = new Duration(1_000_000_000, TimeUnit.DAYS);

private DataSize queryMaxScanRawInputBytes = DataSize.succinctDataSize(1000, PETABYTE);

private int requiredWorkers = 1;
private Duration requiredWorkersMaxWait = new Duration(5, TimeUnit.MINUTES);

Expand Down Expand Up @@ -386,6 +391,18 @@ public QueryManagerConfig setQueryMaxCpuTime(Duration queryMaxCpuTime)
return this;
}

public DataSize getQueryMaxScanRawInputBytes()
{
return this.queryMaxScanRawInputBytes;
}

@Config("query.max-scan-raw-input-bytes")
public QueryManagerConfig setQueryMaxScanRawInputBytes(DataSize queryMaxRawInputBytes)
{
this.queryMaxScanRawInputBytes = queryMaxRawInputBytes;
return this;
}

@Min(1)
public int getRemoteTaskMaxCallbackThreads()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,20 @@ public Duration getTotalCpuTime()
return scheduler.getTotalCpuTime();
}

@Override
public DataSize getRawInputDataSize()
{
SqlQuerySchedulerInterface scheduler = queryScheduler.get();
Optional<QueryInfo> finalQueryInfo = stateMachine.getFinalQueryInfo();
if (finalQueryInfo.isPresent()) {
return finalQueryInfo.get().getQueryStats().getRawInputDataSize();
}
if (scheduler == null) {
return new DataSize(0, BYTE);
}
return scheduler.getRawInputDataSize();
}

@Override
public BasicQueryInfo getBasicQueryInfo()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import com.facebook.airlift.concurrent.ThreadPoolExecutorMBean;
import com.facebook.airlift.log.Logger;
import com.facebook.presto.ExceededCpuLimitException;
import com.facebook.presto.ExceededScanLimitException;
import com.facebook.presto.Session;
import com.facebook.presto.event.QueryMonitor;
import com.facebook.presto.execution.QueryExecution.QueryOutputInfo;
Expand All @@ -29,6 +30,7 @@
import com.facebook.presto.version.EmbedVersion;
import com.google.common.collect.Ordering;
import com.google.common.util.concurrent.ListenableFuture;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;
import org.weakref.jmx.Flatten;
import org.weakref.jmx.Managed;
Expand All @@ -50,6 +52,7 @@

import static com.facebook.airlift.concurrent.Threads.threadsNamed;
import static com.facebook.presto.SystemSessionProperties.getQueryMaxCpuTime;
import static com.facebook.presto.SystemSessionProperties.getQueryMaxScanRawInputBytes;
import static com.facebook.presto.execution.QueryState.RUNNING;
import static com.facebook.presto.spi.StandardErrorCode.GENERIC_INTERNAL_ERROR;
import static com.google.common.collect.ImmutableList.toImmutableList;
Expand All @@ -69,6 +72,7 @@ public class SqlQueryManager
private final QueryTracker<QueryExecution> queryTracker;

private final Duration maxQueryCpuTime;
private final DataSize maxQueryScanPhysicalBytes;

private final ScheduledExecutorService queryManagementExecutor;
private final ThreadPoolExecutorMBean queryManagementExecutorMBean;
Expand All @@ -83,6 +87,7 @@ public SqlQueryManager(ClusterMemoryManager memoryManager, QueryMonitor queryMon
this.embedVersion = requireNonNull(embedVersion, "embedVersion is null");

this.maxQueryCpuTime = queryManagerConfig.getQueryMaxCpuTime();
this.maxQueryScanPhysicalBytes = queryManagerConfig.getQueryMaxScanRawInputBytes();

this.queryManagementExecutor = Executors.newScheduledThreadPool(queryManagerConfig.getQueryManagerExecutorPoolSize(), threadsNamed("query-management-%s"));
this.queryManagementExecutorMBean = new ThreadPoolExecutorMBean((ThreadPoolExecutor) queryManagementExecutor);
Expand All @@ -108,6 +113,12 @@ public void start()
catch (Throwable e) {
log.error(e, "Error enforcing query CPU time limits");
}
try {
enforceScanLimits();
}
catch (Throwable e) {
log.error(e, "Error enforcing query scan bytes limits");
}
}, 1, 1, TimeUnit.SECONDS);
}

Expand Down Expand Up @@ -312,4 +323,19 @@ private void enforceCpuLimits()
}
}
}

/**
* Enforce query scan physical bytes limits
*/
private void enforceScanLimits()
{
for (QueryExecution query : queryTracker.getAllQueries()) {
DataSize rawInputSize = query.getRawInputDataSize();
DataSize sessionlimit = getQueryMaxScanRawInputBytes(query.getSession());
DataSize limit = Ordering.natural().min(maxQueryScanPhysicalBytes, sessionlimit);
if (rawInputSize.compareTo(limit) >= 0) {
query.fail(new ExceededScanLimitException(limit));
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Multimap;
import com.google.common.collect.Sets;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;

import javax.annotation.concurrent.GuardedBy;
Expand Down Expand Up @@ -74,6 +75,7 @@
import static com.google.common.collect.ImmutableList.toImmutableList;
import static com.google.common.collect.Iterables.getOnlyElement;
import static com.google.common.collect.Sets.newConcurrentHashSet;
import static io.airlift.units.DataSize.Unit.BYTE;
import static java.lang.String.format;
import static java.util.Objects.requireNonNull;

Expand Down Expand Up @@ -333,6 +335,17 @@ public synchronized Duration getTotalCpuTime()
return new Duration(millis, TimeUnit.MILLISECONDS);
}

public synchronized DataSize getRawInputDataSize()
{
if (planFragment.getTableScanSchedulingOrder().isEmpty()) {
return new DataSize(0, BYTE);
}
long datasize = getAllTasks().stream()
.mapToLong(task -> task.getTaskInfo().getStats().getRawInputDataSize().toBytes())
.sum();
return DataSize.succinctBytes(datasize);
}

public BasicStageExecutionStats getBasicStageStats()
{
return stateMachine.getBasicStageStats(this::getAllTaskInfo);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.Iterables;
import com.google.common.util.concurrent.ListenableFuture;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;

import java.net.URI;
Expand Down Expand Up @@ -773,6 +774,15 @@ public Duration getTotalCpuTime()
return new Duration(millis, MILLISECONDS);
}

@Override
public DataSize getRawInputDataSize()
{
long datasize = stageExecutions.values().stream()
.mapToLong(stage -> stage.getStageExecution().getRawInputDataSize().toBytes())
.sum();
return DataSize.succinctBytes(datasize);
}

public BasicStageExecutionStats getBasicStageStats()
{
List<BasicStageExecutionStats> stageStats = stageExecutions.values().stream()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
import com.google.common.collect.ImmutableSet;
import com.google.common.collect.ListMultimap;
import com.google.common.util.concurrent.ListenableFuture;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;

import java.net.URI;
Expand Down Expand Up @@ -716,6 +717,16 @@ public Duration getTotalCpuTime()
return new Duration(millis, MILLISECONDS);
}

@Override
public DataSize getRawInputDataSize()
{
long rawInputDataSize = getAllStagesExecutions()
.map(SqlStageExecution::getRawInputDataSize)
.mapToLong(DataSize::toBytes)
.sum();
return DataSize.succinctBytes(rawInputDataSize);
}

@Override
public BasicStageExecutionStats getBasicStageStats()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
import com.facebook.presto.execution.BasicStageExecutionStats;
import com.facebook.presto.execution.StageId;
import com.facebook.presto.execution.StageInfo;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;

public interface SqlQuerySchedulerInterface
Expand All @@ -29,6 +30,8 @@ public interface SqlQuerySchedulerInterface

Duration getTotalCpuTime();

DataSize getRawInputDataSize();

BasicStageExecutionStats getBasicStageStats();

StageInfo getStageInfo();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,16 @@
import com.facebook.airlift.configuration.testing.ConfigAssertions;
import com.facebook.presto.execution.QueryManagerConfig.ExchangeMaterializationStrategy;
import com.google.common.collect.ImmutableMap;
import io.airlift.units.DataSize;
import io.airlift.units.Duration;
import org.testng.annotations.Test;

import java.util.Map;
import java.util.concurrent.TimeUnit;

import static io.airlift.units.DataSize.Unit.MEGABYTE;
import static io.airlift.units.DataSize.Unit.PETABYTE;

public class TestQueryManagerConfig
{
@Test
Expand Down Expand Up @@ -52,6 +56,7 @@ public void testDefaults()
.setQueryMaxRunTime(new Duration(100, TimeUnit.DAYS))
.setQueryMaxExecutionTime(new Duration(100, TimeUnit.DAYS))
.setQueryMaxCpuTime(new Duration(1_000_000_000, TimeUnit.DAYS))
.setQueryMaxScanRawInputBytes(new DataSize(1000, PETABYTE))
.setRequiredWorkers(1)
.setRequiredWorkersMaxWait(new Duration(5, TimeUnit.MINUTES))
.setQuerySubmissionMaxThreads(Runtime.getRuntime().availableProcessors() * 2)
Expand Down Expand Up @@ -86,6 +91,7 @@ public void testExplicitPropertyMappings()
.put("query.max-run-time", "2h")
.put("query.max-execution-time", "3h")
.put("query.max-cpu-time", "2d")
.put("query.max-scan-raw-input-bytes", "1MB")
.put("query.use-streaming-exchange-for-mark-distinct", "true")
.put("query-manager.required-workers", "333")
.put("query-manager.required-workers-max-wait", "33m")
Expand Down Expand Up @@ -117,6 +123,7 @@ public void testExplicitPropertyMappings()
.setQueryMaxRunTime(new Duration(2, TimeUnit.HOURS))
.setQueryMaxExecutionTime(new Duration(3, TimeUnit.HOURS))
.setQueryMaxCpuTime(new Duration(2, TimeUnit.DAYS))
.setQueryMaxScanRawInputBytes(new DataSize(1, MEGABYTE))
.setRequiredWorkers(333)
.setRequiredWorkersMaxWait(new Duration(33, TimeUnit.MINUTES))
.setQuerySubmissionMaxThreads(5)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ public enum StandardErrorCode
EXCEEDED_SPILL_LIMIT(0x0002_0006, INSUFFICIENT_RESOURCES),
EXCEEDED_LOCAL_MEMORY_LIMIT(0x0002_0007, INSUFFICIENT_RESOURCES),
ADMINISTRATIVELY_PREEMPTED(0x0002_0008, INSUFFICIENT_RESOURCES),
EXCEEDED_SCAN_RAW_BYTES_READ_LIMIT(0x0002_0009, INSUFFICIENT_RESOURCES),
/**/;

// Connectors can use error codes starting at the range 0x0100_0000
Expand Down
Loading