Skip to content
1 change: 1 addition & 0 deletions .github/workflows/pinot_compatibility_checks.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths-ignore:
- "contrib/**"
- "docs/**"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pinot_integration_tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths-ignore:
- "contrib/**"
- "docs/**"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pinot_java11_client_compatibility.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths-ignore:
- "contrib/**"
- "docs/**"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pinot_quickstart_tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths-ignore:
- "contrib/**"
- "docs/**"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pinot_unit_tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths-ignore:
- "contrib/**"
- "docs/**"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pinot_vuln_check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ on:
pull_request:
branches:
- master
- 'xiangfu0/codex1/sorted-exchange-ordering'
paths:
- "**/pom.xml"
- "**/package.json"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
Expand Down Expand Up @@ -83,6 +84,7 @@
import org.apache.pinot.query.QueryEnvironment;
import org.apache.pinot.query.mailbox.MailboxService;
import org.apache.pinot.query.planner.explain.AskingServerStageExplainer;
import org.apache.pinot.query.planner.logical.WindowSortAutoPlan;
import org.apache.pinot.query.planner.physical.DispatchablePlanFragment;
import org.apache.pinot.query.planner.physical.DispatchableSubPlan;
import org.apache.pinot.query.planner.plannode.PlanNode;
Expand Down Expand Up @@ -138,6 +140,7 @@ public class MultiStageBrokerRequestHandler extends BaseBrokerRequestHandler {
private final WorkerManager _multiClusterWorkerManager;
private final MailboxService _mailboxService;
private final QueryDispatcher _queryDispatcher;
private final WindowSortAutoTuner _windowSortAutoTuner = new WindowSortAutoTuner();
@Nullable
private final ServerRoutingStatsManager _serverRoutingStatsManager;
private final boolean _explainAskingServerDefault;
Expand Down Expand Up @@ -495,7 +498,16 @@ private QueryEnvironment.CompiledQuery compileQuery(long requestId, String query
Map<String, String> queryOptions = sqlNodeAndOptions.getOptions();

try {
ImmutableQueryEnvironment.Config queryEnvConf = getQueryEnvConf(httpHeaders, queryOptions, requestId);
String database = DatabaseUtils.extractDatabaseFromQueryRequest(queryOptions, httpHeaders);
// The profile key contains no SQL text. Stage/input/collation fingerprints further separate individual
// ordered-window exchanges, even when a query has several windows.
String windowSortMode = queryOptions.getOrDefault(CommonConstants.Broker.Request.QueryOptionKey
.WINDOW_SORT_ON_SENDER, _config.getProperty(CommonConstants.Broker.CONFIG_OF_WINDOW_SORT_ON_SENDER,
CommonConstants.Broker.DEFAULT_WINDOW_SORT_ON_SENDER_MODE));
WindowSortAutoTuner.Session windowSortAutoSession = "auto".equalsIgnoreCase(windowSortMode)
? _windowSortAutoTuner.newSession(Objects.hash(query, database, queryOptions)) : null;
ImmutableQueryEnvironment.Config queryEnvConf = getQueryEnvConf(httpHeaders, queryOptions, requestId,
windowSortAutoSession);
QueryEnvironment queryEnv = new QueryEnvironment(queryEnvConf, _multiClusterRoutingContext);
return callAsync(requestId, query, () -> queryEnv.compile(query, sqlNodeAndOptions), queryTimer);
} catch (WebApplicationException e) {
Expand Down Expand Up @@ -540,6 +552,11 @@ private void checkAuthorization(RequesterIdentity requesterIdentity, RequestCont

private ImmutableQueryEnvironment.Config getQueryEnvConf(HttpHeaders httpHeaders, Map<String, String> queryOptions,
long requestId) {
return getQueryEnvConf(httpHeaders, queryOptions, requestId, null);
}

private ImmutableQueryEnvironment.Config getQueryEnvConf(HttpHeaders httpHeaders, Map<String, String> queryOptions,
long requestId, @Nullable WindowSortAutoPlan windowSortAutoPlan) {
String database = DatabaseUtils.extractDatabaseFromQueryRequest(queryOptions, httpHeaders);
boolean inferPartitionHint = _config.getProperty(CommonConstants.Broker.CONFIG_OF_INFER_PARTITION_HINT,
CommonConstants.Broker.DEFAULT_INFER_PARTITION_HINT);
Expand Down Expand Up @@ -590,9 +607,9 @@ private ImmutableQueryEnvironment.Config getQueryEnvConf(HttpHeaders httpHeaders
int sealedInListThreshold = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_SEALED_IN_LIST_THRESHOLD,
CommonConstants.Broker.DEFAULT_SEALED_IN_LIST_THRESHOLD);
boolean windowSortOnSender = _config.getProperty(
String windowSortOnSenderMode = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_WINDOW_SORT_ON_SENDER,
CommonConstants.Broker.DEFAULT_WINDOW_SORT_ON_SENDER);
CommonConstants.Broker.DEFAULT_WINDOW_SORT_ON_SENDER_MODE);
boolean defaultUnnestColumnPruning = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_UNNEST_COLUMN_PRUNING,
CommonConstants.Broker.DEFAULT_UNNEST_COLUMN_PRUNING);
Expand Down Expand Up @@ -632,7 +649,8 @@ private ImmutableQueryEnvironment.Config getQueryEnvConf(HttpHeaders httpHeaders
.defaultDisabledPlannerRules(_defaultDisabledPlannerRules)
.defaultSortExchangeCopyLimit(sortExchangeCopyThreshold)
.defaultSealedInListThreshold(sealedInListThreshold)
.defaultWindowSortOnSender(windowSortOnSender)
.defaultWindowSortOnSenderMode(windowSortOnSenderMode)
.windowSortAutoPlan(windowSortAutoPlan)
.isKWayMergeSupported(!QueryOptionsUtils.isMultiClusterRoutingEnabled(queryOptions, false)
&& _kWayMergeSupported.getAsBoolean())
.build();
Expand Down Expand Up @@ -828,6 +846,7 @@ private BrokerResponse query(QueryEnvironment.CompiledQuery query, long requestI
}
}

long dispatchToReduceTimeNs = System.nanoTime() - executionStartTimeNs;
BrokerResponseNativeV2 brokerResponse = new BrokerResponseNativeV2();

QueryProcessingException processingException = queryResults.getProcessingException();
Expand All @@ -844,9 +863,8 @@ private BrokerResponse query(QueryEnvironment.CompiledQuery query, long requestI
requestContext.setErrorCode(errorCode);
} else {
brokerResponse.setResultTable(queryResults.getResultTable());
long executionEndTimeNs = System.nanoTime();
updatePhaseTimingForTables(tableNames, BrokerQueryPhase.QUERY_EXECUTION,
executionEndTimeNs - executionStartTimeNs);
dispatchToReduceTimeNs);
}

brokerResponse.setClientRequestId(clientRequestId);
Expand Down Expand Up @@ -885,6 +903,15 @@ private BrokerResponse query(QueryEnvironment.CompiledQuery query, long requestI
}
}

List<QueryDispatcher.QueryResult.StageCoverage> stageCoverage = queryResults.getStageCoverage();
if (processingException == null && brokerResponse.getExceptions().isEmpty()
&& (stageCoverage == null || stageCoverage.stream().noneMatch(c -> c != null
&& (c.getMissing() > 0 || c.getMergeFailed() > 0)))
&& query.getWindowSortAutoPlan() instanceof WindowSortAutoTuner.Session) {
((WindowSortAutoTuner.Session) query.getWindowSortAutoPlan()).observe(queryResults.getQueryStats(),
dispatchToReduceTimeNs);
}

fillOldBrokerResponseStats(brokerResponse, queryResults.getQueryStats(), dispatchableSubPlan,
queryResults.getStageCoverage(), queryResults.getStageStatsTrees());

Expand Down
Loading
Loading