diff --git a/docs/querying/query-context-reference.md b/docs/querying/query-context-reference.md index c485c0231c06..911fed68ac2f 100644 --- a/docs/querying/query-context-reference.md +++ b/docs/querying/query-context-reference.md @@ -71,6 +71,7 @@ Unless otherwise noted, the following parameters apply to all query types, and t |`setProcessingThreadNames`|`true`| Whether processing thread names will be set to `queryType_dataSource_intervals` while processing a query. This aids in interpreting thread dumps, and is on by default. Query overhead can be reduced slightly by setting this to `false`. This has a tiny effect in most scenarios, but can be meaningful in high-QPS, low-per-segment-processing-time scenarios. | |`sqlPlannerBloat`|`1000`|Calcite parameter which controls whether to merge two Project operators when inlining expressions causes complexity to increase. Implemented as a workaround to exception `There are not enough rules to produce a node with desired properties: convention=DRUID, sort=[]` thrown after rejecting the merge of two projects.| |`cloneQueryMode`|`excludeClones`| Indicates whether clone Historicals should be queried by brokers. Clone servers are created by the `cloneServers` Coordinator dynamic configuration. Possible values are `excludeClones`, `includeClones` and `preferClones`. `excludeClones` means that clone Historicals are not queried by the broker. `preferClones` indicates that when given a choice between the clone Historical and the original Historical which is being cloned, the broker chooses the clones. Historicals which are not involved in the cloning process will still be queried. `includeClones` means that broker queries any Historical without regarding clone status. This parameter only affects native queries. MSQ does not query Historicals directly.| +|`queryableHistoricalTiers`|`null`|Set of Historical tier names that may be queried. When set, the Broker only queries Historical servers whose `druid.server.tier` is in this set. Segments without a replica on one of the listed tiers are skipped.| |`realtimeSegmentsMode` |`include`| Controls whether realtime segments are queried. `include` queries all segments, including realtime. `exclude` skips realtime segments. `exclusive` queries only realtime segments. | |`realtimeSegmentsOnly` |`false`| **Deprecated.** Use `realtimeSegmentsMode=exclusive` instead. When set to `true`, this is equivalent to `realtimeSegmentsMode=exclusive`. When set to `false`, this is equivalent to `realtimeSegmentsMode=include`.| @@ -140,4 +141,3 @@ For more information, see the following topics: - [Set query context](./query-context.md) to learn how to configure query context parameters. - [SQL query context](sql-query-context.md) for query context parameters specific to Druid SQL. - [SQL-based ingestion reference](../multi-stage-query/reference/#context-parameters) for context parameters used in SQL-based ingestion (MSQ). - diff --git a/processing/src/main/java/org/apache/druid/query/QueryContext.java b/processing/src/main/java/org/apache/druid/query/QueryContext.java index c29296260023..5023a33de490 100644 --- a/processing/src/main/java/org/apache/druid/query/QueryContext.java +++ b/processing/src/main/java/org/apache/druid/query/QueryContext.java @@ -38,6 +38,7 @@ import java.util.Collections; import java.util.Map; import java.util.Objects; +import java.util.Set; import java.util.TreeMap; /** @@ -237,6 +238,19 @@ public > E getEnum(String key, Class clazz, E defaultValue) return QueryContexts.getAsEnum(key, get(key), clazz, defaultValue); } + /** + * Return a value as a {@link Set} of strings, returning {@code null} if the + * context value is not set. The context value may be a JSON array-like + * {@link java.util.Collection} or a single string. + * + * @throws BadQueryContextException for an invalid value + */ + @Nullable + public Set getStringSet(final String key) + { + return QueryContexts.getAsStringSet(key, get(key)); + } + public Granularity getGranularity(String key, ObjectMapper jsonMapper) { final String granularityString = getString(key); @@ -621,6 +635,12 @@ public CloneQueryMode getCloneQueryMode() ); } + @Nullable + public Set getQueryableHistoricalTiers() + { + return getStringSet(QueryContexts.QUERYABLE_HISTORICAL_TIERS); + } + public boolean getEnableRewriteJoinToFilter() { return getBoolean( diff --git a/processing/src/main/java/org/apache/druid/query/QueryContexts.java b/processing/src/main/java/org/apache/druid/query/QueryContexts.java index 44dffc9a427f..da810ed1f86e 100644 --- a/processing/src/main/java/org/apache/druid/query/QueryContexts.java +++ b/processing/src/main/java/org/apache/druid/query/QueryContexts.java @@ -32,9 +32,13 @@ import javax.annotation.Nullable; import java.math.BigDecimal; import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashSet; import java.util.Map; import java.util.Map.Entry; +import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @@ -67,6 +71,7 @@ public class QueryContexts public static final String MAX_NUMERIC_IN_FILTERS = "maxNumericInFilters"; public static final String CURSOR_AUTO_ARRANGE_FILTERS = "cursorAutoArrangeFilters"; public static final String CLONE_QUERY_MODE = "cloneQueryMode"; + public static final String QUERYABLE_HISTORICAL_TIERS = "queryableHistoricalTiers"; /** * This flag controls whether {@link AggregatorFactory#optimizeForSegment(PerSegmentQueryOptimizationContext)} * is used. It is undocumented because its main purpose is to help developers debug issues with the optimizations. @@ -325,6 +330,30 @@ public static String getAsString( throw badTypeException(key, "a String", value); } + @Nullable + public static Set getAsStringSet( + final String key, + final Object value + ) + { + if (value == null) { + return null; + } else if (value instanceof String) { + return Collections.singleton((String) value); + } else if (value instanceof Collection) { + final Set values = new LinkedHashSet<>(); + for (final Object element : (Collection) value) { + if (!(element instanceof String)) { + throw badValueException(key, "a collection of Strings", value); + } + values.add((String) element); + } + return Collections.unmodifiableSet(values); + } + + throw badTypeException(key, "a String or collection of Strings", value); + } + @Nullable public static Boolean getAsBoolean( final String key, diff --git a/processing/src/test/java/org/apache/druid/query/QueryContextTest.java b/processing/src/test/java/org/apache/druid/query/QueryContextTest.java index d5550bc28dcf..07266cd291dd 100644 --- a/processing/src/test/java/org/apache/druid/query/QueryContextTest.java +++ b/processing/src/test/java/org/apache/druid/query/QueryContextTest.java @@ -23,7 +23,9 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.exc.MismatchedInputException; +import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; import com.google.common.collect.Ordering; import nl.jqno.equalsverifier.EqualsVerifier; import nl.jqno.equalsverifier.Warning; @@ -46,6 +48,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -121,6 +124,37 @@ public void testGetString() assertThrows(BadQueryContextException.class, () -> context.getString("key2")); } + @Test + public void testGetStringSet() + { + final QueryContext context = QueryContext.of( + ImmutableMap.of( + "key1", ImmutableList.of("hot", "cold", "hot"), + "key2", "hot", + "key3", ImmutableList.of("hot", 1), + "key4", 1 + ) + ); + + assertEquals(ImmutableSet.of("hot", "cold"), context.getStringSet("key1")); + assertEquals(Collections.singleton("hot"), context.getStringSet("key2")); + assertNull(context.getStringSet("non-exist")); + + assertThrows(BadQueryContextException.class, () -> context.getStringSet("key3")); + assertThrows(BadQueryContextException.class, () -> context.getStringSet("key4")); + } + + @Test + public void testGetQueryableHistoricalTiers() + { + final QueryContext context = QueryContext.of( + ImmutableMap.of(QueryContexts.QUERYABLE_HISTORICAL_TIERS, ImmutableList.of("hot", "cold")) + ); + + final Set queryableHistoricalTiers = context.getQueryableHistoricalTiers(); + assertEquals(ImmutableSet.of("hot", "cold"), queryableHistoricalTiers); + } + @Test public void testGetBoolean() { diff --git a/server/src/main/java/org/apache/druid/client/CachingClusteredClient.java b/server/src/main/java/org/apache/druid/client/CachingClusteredClient.java index 9305b1b88e9f..a282b96941ea 100644 --- a/server/src/main/java/org/apache/druid/client/CachingClusteredClient.java +++ b/server/src/main/java/org/apache/druid/client/CachingClusteredClient.java @@ -346,13 +346,19 @@ ClusterQueryResult run( final Set segmentServers = computeSegmentsToQuery(timeline, specificSegments); final CloneQueryMode cloneQueryMode = query.context().getCloneQueryMode(); + final Set queryableHistoricalTiers = query.context().getQueryableHistoricalTiers(); @Nullable final byte[] queryCacheKey = cacheKeyManager.computeSegmentLevelQueryCacheKey(); @Nullable final String prevEtag = (String) query.getContext().get(QueryResource.HEADER_IF_NONE_MATCH); if (prevEtag != null) { @Nullable - final String currentEtag = cacheKeyManager.computeResultLevelCachingEtag(segmentServers, cloneQueryMode, queryCacheKey); + final String currentEtag = cacheKeyManager.computeResultLevelCachingEtag( + segmentServers, + cloneQueryMode, + queryableHistoricalTiers, + queryCacheKey + ); if (null != currentEtag) { responseContext.putEntityTag(currentEtag); } @@ -371,7 +377,8 @@ ClusterQueryResult run( final SortedMap> segmentsByServer = groupSegmentsByServer( segmentServers, - cloneQueryMode + cloneQueryMode, + queryableHistoricalTiers ); LazySequence mergedResultSequence = new LazySequence<>(() -> { List> sequencesByInterval = new ArrayList<>(alreadyCachedResults.size() + segmentsByServer.size()); @@ -445,7 +452,10 @@ private Set computeSegmentsToQuery( final Set segments = new LinkedHashSet<>(); final SegmentPruner segmentPruner = ev.getSegmentPruner(); - RealtimeSegmentsMode realtimeSegmentsMode = query.context().getRealtimeSegmentsMode(); + final QueryContext queryContext = query.context(); + final RealtimeSegmentsMode realtimeSegmentsMode = queryContext.getRealtimeSegmentsMode(); + final CloneQueryMode cloneQueryMode = queryContext.getCloneQueryMode(); + final Set queryableHistoricalTiers = queryContext.getQueryableHistoricalTiers(); // Filter unneeded chunks based on partition dimension for (TimelineObjectHolder holder : serversLookup) { final Collection> filteredChunks; @@ -458,7 +468,7 @@ private Set computeSegmentsToQuery( filteredChunks = Sets.newLinkedHashSet(holder.getObject()); } for (PartitionChunk chunk : filteredChunks) { - ServerSelector server = chunk.getObject(); + final ServerSelector server = chunk.getObject(); switch (realtimeSegmentsMode) { case EXCLUSIVE: if (!server.isRealtimeSegment()) { @@ -473,6 +483,10 @@ private Set computeSegmentsToQuery( case INCLUDE: break; } + if (queryableHistoricalTiers != null + && !server.hasQueryableServer(queryableHistoricalTiers, cloneQueryMode)) { + continue; + } final SegmentDescriptor segment = new SegmentDescriptor( holder.getInterval(), holder.getVersion(), @@ -601,13 +615,18 @@ private Cache.NamedKey getCachePopulatorKey(String segmentId, Interval segmentIn private SortedMap> groupSegmentsByServer( Set segments, - CloneQueryMode cloneQueryMode + CloneQueryMode cloneQueryMode, + @Nullable Set queryableHistoricalTiers ) { final SortedMap> serverSegments = new TreeMap<>(); for (SegmentServerSelector segmentServer : segments) { final QueryableDruidServer queryableDruidServer = segmentServer.getServer() - .pick(query, cloneQueryMode); + .pick( + query, + cloneQueryMode, + queryableHistoricalTiers + ); if (queryableDruidServer == null) { log.makeAlert( @@ -820,13 +839,18 @@ byte[] computeSegmentLevelQueryCacheKey() String computeResultLevelCachingEtag( final Set segments, final CloneQueryMode cloneQueryMode, + @Nullable final Set queryableHistoricalTiers, @Nullable byte[] queryCacheKey ) { - Hasher hasher = Hashing.sha1().newHasher(); + final Hasher hasher = Hashing.sha1().newHasher(); boolean hasOnlyHistoricalSegments = true; for (SegmentServerSelector p : segments) { - QueryableDruidServer queryableServer = p.getServer().pick(query, cloneQueryMode); + final QueryableDruidServer queryableServer = p.getServer().pick( + query, + cloneQueryMode, + queryableHistoricalTiers + ); if (queryableServer == null || !queryableServer.getServer().isSegmentReplicationTarget()) { hasOnlyHistoricalSegments = false; break; diff --git a/server/src/main/java/org/apache/druid/client/selector/ServerSelector.java b/server/src/main/java/org/apache/druid/client/selector/ServerSelector.java index 54bf3b2e0fb7..9ca37acdd829 100644 --- a/server/src/main/java/org/apache/druid/client/selector/ServerSelector.java +++ b/server/src/main/java/org/apache/druid/client/selector/ServerSelector.java @@ -21,11 +21,13 @@ import com.google.common.annotations.VisibleForTesting; import com.google.errorprone.annotations.concurrent.GuardedBy; +import it.unimi.dsi.fastutil.ints.Int2ObjectMap; import it.unimi.dsi.fastutil.ints.Int2ObjectRBTreeMap; import org.apache.druid.client.DataSegmentInterner; import org.apache.druid.client.QueryableDruidServer; import org.apache.druid.query.CloneQueryMode; import org.apache.druid.query.Query; +import org.apache.druid.query.QueryContext; import org.apache.druid.server.coordination.DruidServerMetadata; import org.apache.druid.server.coordination.ServerType; import org.apache.druid.timeline.DataSegment; @@ -190,16 +192,123 @@ public List getAllServers(CloneQueryMode cloneQueryMode) } @Nullable - public QueryableDruidServer pick(@Nullable Query query, CloneQueryMode cloneQueryMode) + public QueryableDruidServer pick(@Nullable final Query query, final CloneQueryMode cloneQueryMode) + { + return pick(query, cloneQueryMode, getQueryableHistoricalTiers(query)); + } + + @Nullable + public QueryableDruidServer pick( + @Nullable final Query query, + final CloneQueryMode cloneQueryMode, + @Nullable final Set queryableHistoricalTiers + ) { synchronized (this) { - if (!historicalServers.isEmpty()) { - return historicalTierStrategy.pick(query, filter.getQueryableServers(historicalServers, cloneQueryMode), segment.get()); + if (queryableHistoricalTiers != null && queryableHistoricalTiers.isEmpty()) { + return null; + } + if (!historicalServers.isEmpty() || queryableHistoricalTiers != null) { + final Int2ObjectRBTreeMap> queryableHistoricalServers = + getQueryableHistoricalServers( + filter.getQueryableServers(historicalServers, cloneQueryMode), + queryableHistoricalTiers + ); + return historicalTierStrategy.pick(query, queryableHistoricalServers, segment.get()); } return realtimeTierStrategy.pick(query, realtimeServers, segment.get()); } } + public boolean hasQueryableServer(@Nullable final Query query, final CloneQueryMode cloneQueryMode) + { + return hasQueryableServer(getQueryableHistoricalTiers(query), cloneQueryMode); + } + + public boolean hasQueryableServer( + @Nullable final Set queryableHistoricalTiers, + final CloneQueryMode cloneQueryMode + ) + { + synchronized (this) { + if (queryableHistoricalTiers != null) { + return !queryableHistoricalTiers.isEmpty() + && !historicalServers.isEmpty() + && hasQueryableHistoricalServer( + filter.getQueryableServers(historicalServers, cloneQueryMode), + queryableHistoricalTiers + ); + } + if (!historicalServers.isEmpty()) { + return hasServers(filter.getQueryableServers(historicalServers, cloneQueryMode)); + } + return hasServers(realtimeServers); + } + } + + @Nullable + private static Set getQueryableHistoricalTiers(@Nullable final Query query) + { + if (query == null) { + return null; + } + + final QueryContext queryContext = query.context(); + return queryContext == null ? null : queryContext.getQueryableHistoricalTiers(); + } + + private Int2ObjectRBTreeMap> getQueryableHistoricalServers( + final Int2ObjectRBTreeMap> queryableServers, + @Nullable final Set queryableHistoricalTiers + ) + { + if (queryableHistoricalTiers == null) { + return queryableServers; + } + + final Int2ObjectRBTreeMap> filteredServers = + new Int2ObjectRBTreeMap<>(historicalTierStrategy.getComparator()); + for (final Int2ObjectMap.Entry> entry : queryableServers.int2ObjectEntrySet()) { + final Set priorityServers = new HashSet<>(); + for (final QueryableDruidServer server : entry.getValue()) { + if (queryableHistoricalTiers.contains(server.getServer().getTier())) { + priorityServers.add(server); + } + } + if (!priorityServers.isEmpty()) { + filteredServers.put(entry.getIntKey(), priorityServers); + } + } + + return filteredServers; + } + + private static boolean hasQueryableHistoricalServer( + final Int2ObjectRBTreeMap> queryableServers, + final Set queryableHistoricalTiers + ) + { + for (final Set priorityServers : queryableServers.values()) { + for (final QueryableDruidServer server : priorityServers) { + if (queryableHistoricalTiers.contains(server.getServer().getTier())) { + return true; + } + } + } + + return false; + } + + private static boolean hasServers(final Int2ObjectRBTreeMap> servers) + { + for (final Set priorityServers : servers.values()) { + if (!priorityServers.isEmpty()) { + return true; + } + } + return false; + } + @Override public boolean overshadows(ServerSelector other) { diff --git a/server/src/test/java/org/apache/druid/client/CachingClusteredClientCacheKeyManagerTest.java b/server/src/test/java/org/apache/druid/client/CachingClusteredClientCacheKeyManagerTest.java index 9e043d175d80..b0095dba0795 100644 --- a/server/src/test/java/org/apache/druid/client/CachingClusteredClientCacheKeyManagerTest.java +++ b/server/src/test/java/org/apache/druid/client/CachingClusteredClientCacheKeyManagerTest.java @@ -87,7 +87,7 @@ public void testComputeEtag_nonHistorical() makeHistoricalServerSelector(0), makeRealtimeServerSelector(1) ); - String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, QUERY_CACHE_KEY); + String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, QUERY_CACHE_KEY); Assert.assertNull(actual); } @@ -100,14 +100,14 @@ public void testComputeEtag_DifferentHistoricals() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, QUERY_CACHE_KEY); + String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, QUERY_CACHE_KEY); Assert.assertNotNull(actual1); selectors = ImmutableSet.of( makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, QUERY_CACHE_KEY); + String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, QUERY_CACHE_KEY); Assert.assertNotNull(actual2); Assert.assertEquals("cache key should not change for same server selectors", actual1, actual2); @@ -115,7 +115,7 @@ public void testComputeEtag_DifferentHistoricals() makeHistoricalServerSelector(2), makeHistoricalServerSelector(1) ); - String actual3 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, QUERY_CACHE_KEY); + String actual3 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, QUERY_CACHE_KEY); Assert.assertNotNull(actual3); Assert.assertNotEquals(actual1, actual3); } @@ -129,10 +129,10 @@ public void testComputeEtag_DifferentQueryCacheKey() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, new byte[]{1, 2}); + String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, new byte[]{1, 2}); Assert.assertNotNull(actual1); - String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, new byte[]{3, 4}); + String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, new byte[]{3, 4}); Assert.assertNotNull(actual2); Assert.assertNotEquals(actual1, actual2); } @@ -147,14 +147,14 @@ public void testComputeEtag_nonJoinDataSource() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, FULL_QUERY_CACHE_KEY); + String actual1 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, FULL_QUERY_CACHE_KEY); Assert.assertNotNull(actual1); selectors = ImmutableSet.of( makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null); + String actual2 = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, null); Assert.assertNotNull(actual2); Assert.assertEquals(actual1, actual2); } @@ -170,7 +170,7 @@ public void testComputeEtag_joinWithUnsupportedCaching() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null); + String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, null); Assert.assertNull(actual); } @@ -188,7 +188,7 @@ public void testComputeEtag_noEffectifBySegment() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null); + String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, null); Assert.assertNotNull(actual); } @@ -207,7 +207,7 @@ public void testComputeEtag_noEffectIfUseAndPopulateFalse() makeHistoricalServerSelector(1), makeHistoricalServerSelector(1) ); - String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null); + String actual = keyManager.computeResultLevelCachingEtag(selectors, CloneQueryMode.EXCLUDECLONES, null, null); Assert.assertNotNull(actual); } @@ -295,7 +295,7 @@ private SegmentServerSelector makeServerSelector(boolean isHistorical, int parti SegmentId segmentId = SegmentId.dummy("data-source", partitionNumber); DataSegment segment = DataSegment.builder(segmentId).shardSpec(new NumberedShardSpec(partitionNumber, 10)).build(); expect(server.isSegmentReplicationTarget()).andReturn(isHistorical).anyTimes(); - expect(serverSelector.pick(query, CloneQueryMode.EXCLUDECLONES)).andReturn(queryableDruidServer).anyTimes(); + expect(serverSelector.pick(query, CloneQueryMode.EXCLUDECLONES, null)).andReturn(queryableDruidServer).anyTimes(); expect(queryableDruidServer.getServer()).andReturn(server).anyTimes(); expect(serverSelector.getSegment()).andReturn(segment).anyTimes(); replay(serverSelector, queryableDruidServer, server); diff --git a/server/src/test/java/org/apache/druid/client/CachingClusteredClientTest.java b/server/src/test/java/org/apache/druid/client/CachingClusteredClientTest.java index 94c259314dea..2952330a72a0 100644 --- a/server/src/test/java/org/apache/druid/client/CachingClusteredClientTest.java +++ b/server/src/test/java/org/apache/druid/client/CachingClusteredClientTest.java @@ -3222,6 +3222,61 @@ public void testRealtimeSegmentsModeExclude() Assert.assertEquals(1, remainingResponseMap.get(queryInclude.getId()).intValue()); } + @Test + public void testQueryableHistoricalTiersQueryContext() + { + final Interval interval = Intervals.of("2016-01-01/2016-01-02"); + final Interval queryInterval = Intervals.of("2016-01-01T14:00:00/2016-01-02T14:00:00"); + final DataSegment dataSegment = DataSegment.builder( + SegmentId.of("dataSource", interval, "ver", NoneShardSpec.instance()) + ) + .loadSpec(ImmutableMap.of("type", "hdfs", "path", "/tmp")) + .dimensions(ImmutableList.of("product")) + .metrics(ImmutableList.of("visited_sum")) + .binaryVersion(9) + .size(12334) + .build(); + final ServerSelector selector = new ServerSelector( + dataSegment, + new HighestPriorityTierSelectorStrategy(new RandomServerSelectorStrategy()), + HistoricalFilter.IDENTITY_FILTER + ); + selector.addServerAndUpdateSegment(new QueryableDruidServer(servers[0], null), dataSegment); + timeline.add(interval, "ver", new SingleElementPartitionChunk<>(selector)); + + final TimeBoundaryQuery disallowedTierQuery = Druids.newTimeBoundaryQueryBuilder() + .dataSource(DATA_SOURCE) + .intervals(new MultipleIntervalSegmentSpec( + ImmutableList.of(queryInterval) + )) + .context(ImmutableMap.of( + QueryContexts.QUERYABLE_HISTORICAL_TIERS, + ImmutableList.of("hot") + )) + .randomQueryId() + .build(); + final TimeBoundaryQuery allowedTierQuery = Druids.newTimeBoundaryQueryBuilder() + .dataSource(DATA_SOURCE) + .intervals(new MultipleIntervalSegmentSpec( + ImmutableList.of(queryInterval) + )) + .context(ImmutableMap.of( + QueryContexts.QUERYABLE_HISTORICAL_TIERS, + ImmutableList.of("bye") + )) + .randomQueryId() + .build(); + + final ResponseContext responseContext = initializeResponseContext(); + getDefaultQueryRunner().run(QueryPlus.wrap(disallowedTierQuery), responseContext); + getDefaultQueryRunner().run(QueryPlus.wrap(allowedTierQuery), responseContext); + + final Map remainingResponseMap = + (Map) responseContext.get(ResponseContext.Keys.REMAINING_RESPONSES_FROM_QUERY_SERVERS); + Assert.assertEquals(0, remainingResponseMap.get(disallowedTierQuery.getId()).intValue()); + Assert.assertEquals(1, remainingResponseMap.get(allowedTierQuery.getId()).intValue()); + } + @SuppressWarnings("unchecked") private QueryRunner getDefaultQueryRunner() { diff --git a/server/src/test/java/org/apache/druid/client/selector/ServerSelectorTest.java b/server/src/test/java/org/apache/druid/client/selector/ServerSelectorTest.java index 015fb3dfa4e4..7353d7fdb969 100644 --- a/server/src/test/java/org/apache/druid/client/selector/ServerSelectorTest.java +++ b/server/src/test/java/org/apache/druid/client/selector/ServerSelectorTest.java @@ -25,8 +25,13 @@ import org.apache.druid.client.DruidServer; import org.apache.druid.client.QueryableDruidServer; import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.query.CloneQueryMode; +import org.apache.druid.query.Druids; +import org.apache.druid.query.Query; +import org.apache.druid.query.QueryContexts; import org.apache.druid.server.coordination.ServerType; import org.apache.druid.timeline.DataSegment; +import org.apache.druid.timeline.SegmentId; import org.apache.druid.timeline.partition.NoneShardSpec; import org.apache.druid.timeline.partition.TombstoneShardSpec; import org.easymock.EasyMock; @@ -168,4 +173,116 @@ public void testSegmentWithData() Assert.assertTrue(selector.hasData()); } + @Test + public void testQueryableHistoricalTiersFiltersHistoricalServers() + { + final DataSegment segment = createSegmentForQueryableHistoricalTiers(); + final ServerSelector selector = new ServerSelector( + segment, + new HighestPriorityTierSelectorStrategy(new RandomServerSelectorStrategy()), + HistoricalFilter.IDENTITY_FILTER + ); + final DruidServer coldServer = new DruidServer( + "cold", + "cold", + null, + 0, + null, + ServerType.HISTORICAL, + "cold", + 10 + ); + final DruidServer hotServer = new DruidServer( + "hot", + "hot", + null, + 0, + null, + ServerType.HISTORICAL, + "hot", + 0 + ); + selector.addServerAndUpdateSegment( + new QueryableDruidServer(coldServer, EasyMock.createMock(DirectDruidClient.class)), + segment + ); + selector.addServerAndUpdateSegment( + new QueryableDruidServer(hotServer, EasyMock.createMock(DirectDruidClient.class)), + segment + ); + + final Query hotTierQuery = Druids.newTimeBoundaryQueryBuilder() + .dataSource("test") + .intervals("2012/2013") + .context(ImmutableMap.of( + QueryContexts.QUERYABLE_HISTORICAL_TIERS, + ImmutableList.of("hot") + )) + .build(); + final Query missingTierQuery = Druids.newTimeBoundaryQueryBuilder() + .dataSource("test") + .intervals("2012/2013") + .context(ImmutableMap.of( + QueryContexts.QUERYABLE_HISTORICAL_TIERS, + ImmutableList.of("warm") + )) + .build(); + + Assert.assertEquals( + hotServer, + selector.pick(hotTierQuery, CloneQueryMode.EXCLUDECLONES).getServer() + ); + Assert.assertTrue(selector.hasQueryableServer(hotTierQuery, CloneQueryMode.EXCLUDECLONES)); + Assert.assertNull(selector.pick(missingTierQuery, CloneQueryMode.EXCLUDECLONES)); + Assert.assertFalse(selector.hasQueryableServer(missingTierQuery, CloneQueryMode.EXCLUDECLONES)); + } + + @Test + public void testQueryableHistoricalTiersDoesNotFallBackToRealtime() + { + final DataSegment segment = createSegmentForQueryableHistoricalTiers(); + final ServerSelector selector = new ServerSelector( + segment, + new HighestPriorityTierSelectorStrategy(new RandomServerSelectorStrategy()), + HistoricalFilter.IDENTITY_FILTER + ); + selector.addServerAndUpdateSegment( + new QueryableDruidServer( + new DruidServer("realtime", "realtime", null, 0, null, ServerType.REALTIME, "hot", 0), + EasyMock.createMock(DirectDruidClient.class) + ), + segment + ); + + final Query hotTierQuery = Druids.newTimeBoundaryQueryBuilder() + .dataSource("test") + .intervals("2012/2013") + .context(ImmutableMap.of( + QueryContexts.QUERYABLE_HISTORICAL_TIERS, + ImmutableList.of("hot") + )) + .build(); + + Assert.assertNull(selector.pick(hotTierQuery, CloneQueryMode.EXCLUDECLONES)); + Assert.assertFalse(selector.hasQueryableServer(hotTierQuery, CloneQueryMode.EXCLUDECLONES)); + } + + private static DataSegment createSegmentForQueryableHistoricalTiers() + { + return DataSegment.builder( + SegmentId.of( + "test_broker_server_view", + Intervals.of("2012/2013"), + "v1", + NoneShardSpec.instance() + ) + ) + .loadSpec(ImmutableMap.of("type", "local", "path", "somewhere")) + .dimensions(ImmutableList.of()) + .metrics(ImmutableList.of()) + .binaryVersion(9) + .size(0) + .build(); + } + }