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 @@ -22,6 +22,7 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.inject.Injector;
import org.apache.druid.client.TimelineServerView;
import org.apache.druid.discovery.DruidNodeDiscovery;
import org.apache.druid.error.DruidException;
import org.apache.druid.indexing.common.TaskLockType;
import org.apache.druid.indexing.common.actions.TaskActionClient;
Expand Down Expand Up @@ -52,6 +53,7 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;

/**
Expand Down Expand Up @@ -91,6 +93,7 @@ public class DartControllerContext implements ControllerContext
private final List<InputSpecSlicerProvider> inputSpecSlicerProviders;
private final ServiceEmitter emitter;
private final QueryContext context;
private final DruidNodeDiscovery dartWorkerDiscovery;

public DartControllerContext(
final Injector injector,
Expand All @@ -101,7 +104,8 @@ public DartControllerContext(
final TimelineServerView serverView,
final List<InputSpecSlicerProvider> inputSpecSlicerProviders,
final ServiceEmitter emitter,
final QueryContext context
final QueryContext context,
final DruidNodeDiscovery dartWorkerDiscovery
)
{
this.injector = injector;
Expand All @@ -113,6 +117,7 @@ public DartControllerContext(
this.inputSpecSlicerProviders = inputSpecSlicerProviders;
this.emitter = emitter;
this.context = context;
this.dartWorkerDiscovery = dartWorkerDiscovery;
}

@Override
Expand All @@ -126,17 +131,38 @@ public ControllerQueryKernelConfig queryKernelConfig(final MSQSpec querySpec)
{
final List<DruidServerMetadata> servers = serverView.getDruidServerMetadatas();

final Set<String> dartWorkerHosts =
dartWorkerDiscovery.getAllNodes()
.stream()
.map(node -> node.getDruidNode().getHostAndPortToUse())
.collect(Collectors.toSet());

// Lock in the list of workers when creating the kernel config. There is a race here: the serverView itself is
// allowed to float. If a segment moves to a new server that isn't part of our list after the WorkerManager is
// created, we won't be able to find a valid server for certain segments. This isn't expected to be a problem,
// since the serverView is referenced shortly after the worker list is created.
final List<String> workerIds = new ArrayList<>(servers.size());
for (final DruidServerMetadata server : servers) {
if (server.getType() == ServerType.HISTORICAL) {
if (server.getType() == ServerType.HISTORICAL && dartWorkerHosts.contains(server.getHost())) {
workerIds.add(WorkerId.fromDruidServerMetadata(server, queryId()).toString());
}
}

// Fail fast rather than running a query with no workers.
if (workerIds.isEmpty()) {
final boolean anyHistoricals = servers.stream().anyMatch(s -> s.getType() == ServerType.HISTORICAL);
throw DruidException.forPersona(DruidException.Persona.OPERATOR)
.ofCategory(DruidException.Category.RUNTIME_FAILURE)
.build(
anyHistoricals
? "No Dart workers are available: Historicals are present but none advertise a Dart "
+ "worker. Set druid.msq.dart.enabled=true on the Historicals that should run Dart "
+ "queries."
: "No Dart workers are available: no Historicals are currently available to run the "
+ "query."
);
}

// Shuffle workerIds, so we don't bias towards specific servers when running multiple queries concurrently. For any
// given query, lower-numbered workers tend to do more work, because the controller prefers using lower-numbered
// workers when maxWorkerCount for a stage is less than the total number of workers.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,17 @@
import com.google.inject.Inject;
import com.google.inject.Injector;
import org.apache.druid.client.TimelineServerView;
import org.apache.druid.discovery.DruidNodeDiscovery;
import org.apache.druid.discovery.DruidNodeDiscoveryProvider;
import org.apache.druid.discovery.NodeRole;
import org.apache.druid.guice.annotations.EscalatedGlobal;
import org.apache.druid.guice.annotations.Json;
import org.apache.druid.guice.annotations.Self;
import org.apache.druid.guice.annotations.Smile;
import org.apache.druid.java.util.emitter.service.ServiceEmitter;
import org.apache.druid.msq.dart.Dart;
import org.apache.druid.msq.dart.worker.DartWorkerClientImpl;
import org.apache.druid.msq.dart.worker.DartWorkerService;
import org.apache.druid.msq.exec.ControllerContext;
import org.apache.druid.msq.exec.MemoryIntrospector;
import org.apache.druid.msq.input.InputSpecSlicerProvider;
Expand All @@ -52,6 +56,7 @@ public class DartControllerContextFactoryImpl implements DartControllerContextFa
protected final MemoryIntrospector memoryIntrospector;
protected final List<InputSpecSlicerProvider> inputSpecSlicerProviders;
protected final ServiceEmitter emitter;
protected final DruidNodeDiscovery dartWorkerDiscovery;

@Inject
public DartControllerContextFactoryImpl(
Expand All @@ -63,7 +68,8 @@ public DartControllerContextFactoryImpl(
final MemoryIntrospector memoryIntrospector,
final TimelineServerView serverView,
@Dart final Set<InputSpecSlicerProvider> inputSpecSlicerProviders,
final ServiceEmitter emitter
final ServiceEmitter emitter,
final DruidNodeDiscoveryProvider discoveryProvider
)
{
this.injector = injector;
Expand All @@ -75,6 +81,8 @@ public DartControllerContextFactoryImpl(
this.memoryIntrospector = memoryIntrospector;
this.inputSpecSlicerProviders = List.copyOf(inputSpecSlicerProviders);
this.emitter = emitter;
this.dartWorkerDiscovery =
discoveryProvider.getForServiceAndRoles(DartWorkerService.NAME, Set.of(NodeRole.HISTORICAL));
}

@Override
Expand All @@ -90,7 +98,8 @@ public ControllerContext newContext(final QueryContext context)
serverView,
inputSpecSlicerProviders,
emitter,
context
context,
dartWorkerDiscovery
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@
import org.apache.druid.messages.client.MessageRelayFactory;
import org.apache.druid.messages.client.MessageRelays;
import org.apache.druid.msq.dart.controller.messages.ControllerMessage;
import org.apache.druid.msq.dart.worker.DartWorkerService;

import java.util.Set;

/**
* Specialized {@link MessageRelays} for Dart controllers.
Expand All @@ -35,6 +38,11 @@ public DartMessageRelays(
final MessageRelayFactory<ControllerMessage> messageRelayFactory
)
{
super(() -> discoveryProvider.getForNodeRole(NodeRole.HISTORICAL), messageRelayFactory);
// Only relay with Historicals that run a Dart worker (advertise DartWorkerService); Dart-disabled Historicals
// expose no outbox. This is the same discovery handle the controller uses to enroll workers.
super(
() -> discoveryProvider.getForServiceAndRoles(DartWorkerService.NAME, Set.of(NodeRole.HISTORICAL)),
messageRelayFactory
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,10 @@
import com.google.inject.Key;
import com.google.inject.Module;
import com.google.inject.Provides;
import com.google.inject.multibindings.ProvidesIntoSet;
import com.google.inject.name.Named;
import org.apache.druid.discovery.DruidNodeDiscoveryProvider;
import org.apache.druid.discovery.DruidService;
import org.apache.druid.discovery.NodeRole;
import org.apache.druid.guice.Jerseys;
import org.apache.druid.guice.JsonConfigProvider;
Expand Down Expand Up @@ -57,6 +60,7 @@
import org.apache.druid.msq.dart.worker.DartWorkerContextFactory;
import org.apache.druid.msq.dart.worker.DartWorkerContextFactoryImpl;
import org.apache.druid.msq.dart.worker.DartWorkerRunner;
import org.apache.druid.msq.dart.worker.DartWorkerService;
import org.apache.druid.msq.dart.worker.http.DartWorkerResource;
import org.apache.druid.msq.exec.MemoryIntrospector;
import org.apache.druid.msq.guice.MSQBinders;
Expand Down Expand Up @@ -118,6 +122,19 @@ public void configure(Binder binder)
.in(LazySingleton.class);
}

/**
* Advertise {@link DartWorkerService} in node discovery. Contributed from {@link ActualModule}, which is only
* installed when Dart is enabled, so advertisement tracks actually running a Dart worker. Merges into the same
* {@code @Named("historical")} service set announced by
* {@link org.apache.druid.guice.HistoricalServiceModule}.
*/
@ProvidesIntoSet
@Named(NodeRole.HISTORICAL_JSON_NAME)
public Class<? extends DruidService> getDartWorkerService()
{
return DartWorkerService.class;
}

@Provides
@ManageLifecycle
public DartWorkerRunner createWorkerRunner(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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 org.apache.druid.msq.dart.worker;

import org.apache.druid.discovery.DruidService;

/**
* No-payload {@link DruidService} advertised in node discovery by a Historical whose Dart worker runtime is
* installed (i.e. {@code druid.msq.dart.enabled=true}). The Dart controller enrolls workers only when they are
* Historicals that advertise this capability.
*/
public class DartWorkerService extends DruidService
{
public static final String NAME = "dartWorkerService";

@Override
public String getName()
{
return NAME;
}

@Override
public boolean equals(Object o)
{
return o != null && getClass() == o.getClass();
}

@Override
public int hashCode()
{
return DartWorkerService.class.hashCode();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
import org.apache.druid.msq.counters.StorageCounters;
import org.apache.druid.msq.counters.SuperSorterProgressTrackerCounter;
import org.apache.druid.msq.counters.WarningCounters;
import org.apache.druid.msq.dart.worker.DartWorkerService;
import org.apache.druid.msq.indexing.IndexerControllerContextFactory;
import org.apache.druid.msq.indexing.IndexerSegmentsInputSliceReaderProvider;
import org.apache.druid.msq.indexing.IndexerTableInputSpecSlicerProvider;
Expand Down Expand Up @@ -229,6 +230,10 @@ public List<? extends Module> getJacksonModules()

module.registerSubtypes(new NamedType(MSQCompactionRunner.class, MSQCompactionRunner.TYPE));

// Registered here rather than in DartWorkerModule so that every process that discovers Historicals can parse
// their DartWorkerService announcement.
module.registerSubtypes(new NamedType(DartWorkerService.class, DartWorkerService.NAME));

FAULT_CLASSES.forEach(module::registerSubtypes);
module.addSerializer(new CounterSnapshotsSerializer());
return Collections.singletonList(module);
Expand Down
Loading
Loading