Skip to content
Open
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 @@ -25,6 +25,7 @@

import javax.annotation.Nullable;
import java.util.List;
import java.util.Map;

/**
* {@link Cursor} that walks a sequence of per-group cursors back-to-back, presenting them to the caller as a single
Expand All @@ -40,6 +41,7 @@ public final class ConcatenatingCursor implements Cursor
private final List<Supplier<CursorHolder>> holderSuppliers;
private final List<Object[]> clusteringValuesByGroup;
private final ClusteringColumnSelectorFactory wrapperFactory;
private final ColumnSelectorFactory exposedFactory;

private int currentIdx;
@Nullable
Expand All @@ -49,7 +51,8 @@ public final class ConcatenatingCursor implements Cursor
public ConcatenatingCursor(
List<Supplier<CursorHolder>> holderSuppliers,
List<Object[]> clusteringValuesByGroup,
ClusteringColumnSelectorFactory wrapperFactory
ClusteringColumnSelectorFactory wrapperFactory,
Map<String, String> virtualColumnRemap
)
{
if (holderSuppliers.size() != clusteringValuesByGroup.size()) {
Expand All @@ -65,6 +68,9 @@ public ConcatenatingCursor(
this.holderSuppliers = holderSuppliers;
this.clusteringValuesByGroup = clusteringValuesByGroup;
this.wrapperFactory = wrapperFactory;
this.exposedFactory = virtualColumnRemap.isEmpty()
? wrapperFactory
: new RemapColumnSelectorFactory(wrapperFactory, virtualColumnRemap);
this.currentIdx = -1;
}

Expand Down Expand Up @@ -101,7 +107,7 @@ private void advanceToNextNonEmptyGroup()
public ColumnSelectorFactory getColumnSelectorFactory()
{
initializeIfNeeded();
return wrapperFactory;
return exposedFactory;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.druid.segment.vector.ConcatenatingVectorCursor;
import org.apache.druid.segment.vector.MultiValueDimensionVectorSelector;
import org.apache.druid.segment.vector.ReadableVectorInspector;
import org.apache.druid.segment.vector.RemapVectorColumnSelectorFactory;
import org.apache.druid.segment.vector.SingleValueDimensionVectorSelector;
import org.apache.druid.segment.vector.VectorColumnSelectorFactory;
import org.apache.druid.segment.vector.VectorCursor;
Expand Down Expand Up @@ -200,13 +201,43 @@ private CursorHolder makeSingleGroupClusteredCursorHolder(
);
}

// groupIndex exposes the group's clustering columns as constant columns, no selector wrapper is needed
return new QueryableIndexCursorHolder(
groupIndex,
plan.rebuildCursorBuildSpec(spec, valueGroup),
QueryableIndexTimeBoundaryInspector.create(groupIndex),
valueGroup.getSummary().getOrdering()
);
// groupIndex exposes the group's clustering columns as constant columns, so no clustering-specific selector
// wrapper is needed. When the query has virtual columns equivalent to a materialized column, wrap the selector
// factories with Remap*ColumnSelectorFactory so those query VCs resolve to the materialized column instead of
// recomputing (the matched VCs were already dropped from the per-group build spec by rebuildCursorBuildSpec).
final CursorBuildSpec groupSpec = plan.rebuildCursorBuildSpec(spec, valueGroup);
final QueryableIndexTimeBoundaryInspector groupTimeBoundary =
QueryableIndexTimeBoundaryInspector.create(groupIndex);
final List<OrderBy> ordering = valueGroup.getSummary().getOrdering();
if (plan.virtualColumnRemap().isEmpty()) {
return new QueryableIndexCursorHolder(groupIndex, groupSpec, groupTimeBoundary, ordering);
}
return new QueryableIndexCursorHolder(groupIndex, groupSpec, groupTimeBoundary, ordering)
{
@Override
protected ColumnSelectorFactory makeColumnSelectorFactoryForOffset(
ColumnCache columnCache,
Offset baseOffset
)
{
return new RemapColumnSelectorFactory(
super.makeColumnSelectorFactoryForOffset(columnCache, baseOffset),
plan.virtualColumnRemap()
);
}

@Override
protected VectorColumnSelectorFactory makeVectorColumnSelectorFactoryForOffset(
ColumnCache columnCache,
VectorOffset baseOffset
)
{
return new RemapVectorColumnSelectorFactory(
super.makeVectorColumnSelectorFactoryForOffset(columnCache, baseOffset),
plan.virtualColumnRemap()
);
}
};
}

/**
Expand Down Expand Up @@ -271,12 +302,14 @@ private CursorHolder makeMultiGroupClusteredCursorHolder(
final ConcatenatingCursor cursor = new ConcatenatingCursor(
holderSuppliers,
clusteringValuesByGroup,
wrapperFactory
wrapperFactory,
plan.virtualColumnRemap()
);
final ConcatenatingVectorCursor vectorCursor = new ConcatenatingVectorCursor(
holderSuppliers,
clusteringValuesByGroup,
vectorWrapperFactory
vectorWrapperFactory,
plan.virtualColumnRemap()
);

// each group gets a different rewritten filter, so the conservative thing to do here is require the original query
Expand All @@ -285,8 +318,8 @@ private CursorHolder makeMultiGroupClusteredCursorHolder(
final Filter queryFilter = spec.getFilter();
final boolean filterCanVectorize =
queryFilter == null || queryFilter.canVectorizeMatcher(spec.getVirtualColumns().wrapInspector(this));
// we still check that the first holder is vectorizable to make sure all the non-filter parts can be vectorized
final boolean canVectorize = filterCanVectorize && holderSuppliers.get(0).get().canVectorize();
// we still check that the first holder is vectorizable to make sure all the non-filter parts can be vectorized.
final boolean canVectorize = filterCanVectorize && holderSuppliers.getFirst().get().canVectorize();

return new CursorHolder()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,9 +175,14 @@ private CursorHolder makeClusteredCursorHolder(CursorBuildSpec spec, OnHeapClust
final ClusteringColumnSelectorFactory wrapperFactory = new ClusteringColumnSelectorFactory(
ClusteringColumnSelectorFactory.UNINITIALIZED_DELEGATE,
clusteringColumns,
clusteringValuesByGroup.get(0)
clusteringValuesByGroup.getFirst()
);
final ConcatenatingCursor cursor = new ConcatenatingCursor(
holderSuppliers,
clusteringValuesByGroup,
wrapperFactory,
plan.virtualColumnRemap()
);
final ConcatenatingCursor cursor = new ConcatenatingCursor(holderSuppliers, clusteringValuesByGroup, wrapperFactory);

return new CursorHolder()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,17 @@

import org.apache.druid.query.filter.Filter;
import org.apache.druid.segment.CursorBuildSpec;
import org.apache.druid.segment.VirtualColumn;
import org.apache.druid.segment.VirtualColumns;
import org.apache.druid.segment.filter.FalseFilter;
import org.apache.druid.segment.filter.TrueFilter;

import javax.annotation.Nullable;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Function;

/**
Expand All @@ -39,14 +45,17 @@ public final class ClusterGroupQueryPlan
{
private final List<TableClusterGroupSpec> survivingGroups;
private final Function<TableClusterGroupSpec, Filter> rewriter;
private final Map<String, String> virtualColumnRemap;

ClusterGroupQueryPlan(
List<TableClusterGroupSpec> survivingGroups,
Function<TableClusterGroupSpec, Filter> rewriter
Function<TableClusterGroupSpec, Filter> rewriter,
Map<String, String> virtualColumnRemap
)
{
this.survivingGroups = survivingGroups;
this.rewriter = rewriter;
this.virtualColumnRemap = virtualColumnRemap;
}

/**
Expand All @@ -58,6 +67,16 @@ public List<TableClusterGroupSpec> survivingGroups()
return survivingGroups;
}

/**
* Query-level mapping of {@code queryVirtualColumnOutputName -> materializedColumnName} for query virtual columns
* that are equivalent to a clustering column or a non-clustering materialized column in the clustered base table.
* Empty when no query virtual column has a materialized equivalent (or there are no query virtual columns).
*/
public Map<String, String> virtualColumnRemap()
{
return virtualColumnRemap;
}

/**
* Per-group rewrite of the query's filter against {@code group}'s constant clustering tuple: clustering-column
* leaves collapse to {@link TrueFilter} / {@link FalseFilter} and fold through AND / OR / NOT, so the rewritten
Expand All @@ -76,21 +95,60 @@ public Filter rewriteFor(TableClusterGroupSpec group)

/**
* Rebuild {@code spec} for {@code group}'s per-group cursor by swapping in this plan's per-group filter rewrite
* (see {@link #rewriteFor}). Returns {@code spec} unchanged when there is no filter, or when the rewrite is
* identical to the original (no clustering leaves folded), so the common no-op case allocates nothing. Shared by
* the {@link org.apache.druid.segment.QueryableIndexCursorFactory} (historical) and
* {@link org.apache.druid.segment.incremental.IncrementalIndexCursorFactory} (realtime) clustered dispatch.
* (see {@link #rewriteFor}) and, when {@link #virtualColumnRemap()} is non-empty, substituting any query virtual
* columns that are equivalent to a materialized column.
* <p>
* When there is a remap, the matched query virtual columns (the remap keys) are dropped from the spec's virtual
* columns, and the per-group filter's required columns are rewritten via the remap so that non-clustering-VC filter
* leaves point at the materialized physical column (clustering-VC leaves were already folded to TRUE / FALSE by the
* per-group rewrite walk, so they won't reference the dropped virtual columns). The grouping / select / aggregator
* references are served instead by the {@link org.apache.druid.segment.RemapColumnSelectorFactory} the cursor
* factory wraps around the {@link ClusteringColumnSelectorFactory}.
Comment on lines +101 to +106

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

clustering-VC leaves were already folded to TRUE / FALSE by the per-group rewrite walk, so they won't reference the dropped virtual columns

This is a not entirely true. walkClusterGroupFilter only folds Equality, In, and Null. Everything else falls through unchanged. This will result in wrong results when something like a Range filter on a query vc with equivalent clustering column survives the walk.

Repro wrong empty results:

  • tenant_lower := lower(tenant) clustered vc
  • v0 := lower(tenant) ⇒ remap v0 → tenant_lower
  • same test tenants as always
  • RangeFilter("v0", ["a","z"])
  • Both groups survive and you expect all rows but get empty result back

*/
public CursorBuildSpec rebuildCursorBuildSpec(CursorBuildSpec spec, TableClusterGroupSpec group)
{
if (spec.getFilter() == null) {
return spec;
if (virtualColumnRemap.isEmpty()) {
if (spec.getFilter() == null) {
return spec;
}
final Filter rewritten = rewriteFor(group);
if (rewritten == spec.getFilter()) {
return spec;
}
return CursorBuildSpec.builder(spec).setFilter(rewritten).build();
}
final Filter rewritten = rewriteFor(group);
if (rewritten == spec.getFilter()) {
return spec;

// Drop the remapped query virtual columns from the per-group spec; the materialized columns they map to are
// either the per-group constant clustering column or a per-group physical column, both served directly by the
// ClusteringColumnSelectorFactory (the cursor factory additionally wraps it with a RemapColumnSelectorFactory so
// the original query VC names resolve to the materialized columns).
final List<VirtualColumn> prunedVcs = new ArrayList<>();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

naming seems backwards. aren't these actually whare are not being pruned out of the query?

for (VirtualColumn vc : spec.getVirtualColumns().getVirtualColumns()) {
if (!virtualColumnRemap.containsKey(vc.getOutputName())) {
prunedVcs.add(vc);
}
}
return CursorBuildSpec.builder(spec).setFilter(rewritten).build();

// Per-group filter rewrite first (folds clustering leaves), then remap any surviving non-clustering-VC leaves to
// the materialized physical column.
Filter rewritten = spec.getFilter() == null ? null : rewriteFor(group);
if (rewritten != null) {
rewritten = Projections.rewriteFilterRequiredColumns(rewritten, virtualColumnRemap);
}

final CursorBuildSpec.CursorBuildSpecBuilder builder = CursorBuildSpec.builder(spec)
.setFilter(rewritten)
.setVirtualColumns(VirtualColumns.create(prunedVcs));

// The dropped query VCs now read their materialized target columns, so add them
final Set<String> physicalColumns = spec.getPhysicalColumns();
if (physicalColumns != null) {
final Set<String> withTargets = new HashSet<>(physicalColumns);
withTargets.addAll(virtualColumnRemap.values());
builder.setPhysicalColumns(withTargets);
}

return builder.build();
}
}

Loading
Loading