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 @@ -20,31 +20,80 @@
package org.apache.druid.query.operator.window.value;

import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.druid.error.DruidException;
import org.apache.druid.query.operator.window.Processor;
import org.apache.druid.query.operator.window.WindowFrame;
import org.apache.druid.query.rowsandcols.RowsAndColumns;
import org.apache.druid.query.rowsandcols.column.ColumnAccessor;
import org.apache.druid.query.rowsandcols.column.ConstantObjectColumn;
import org.apache.druid.query.rowsandcols.column.ObjectArrayColumn;

import javax.annotation.Nullable;
import java.util.Objects;

public class WindowFirstProcessor extends WindowValueProcessorBase
{
@Nullable
private final WindowFrame frame;

@JsonCreator
public WindowFirstProcessor(
@JsonProperty("inputColumn") String inputColumn,
@JsonProperty("outputColumn") String outputColumn
@JsonProperty("outputColumn") String outputColumn,
@JsonProperty("frame") @Nullable WindowFrame frame
)
{
super(inputColumn, outputColumn);
this.frame = frame;
}

@Nullable
@JsonProperty("frame")
@JsonInclude(JsonInclude.Include.NON_NULL)
public WindowFrame getFrame()
{
return frame;
}

@Override
public RowsAndColumns process(RowsAndColumns incomingPartition)
public RowsAndColumns process(RowsAndColumns input)
{
final int numRows = input.numRows();

if (numRows == 0) {
throw DruidException.defensive("Called with an input partition of size 0. The call site needs to not do that.");
}

if (frame == null) {
return processInternal(
input,
column -> {
final ColumnAccessor accessor = column.toAccessor();
return new ConstantObjectColumn(accessor.getObject(0), accessor.numRows(), accessor.getType());
}
);
}

return processInternal(
incomingPartition,
input,
column -> {
final ColumnAccessor accessor = column.toAccessor();
return new ConstantObjectColumn(accessor.getObject(0), accessor.numRows(), accessor.getType());
final Object[] results = new Object[numRows];
computeFirstOrLastValueFramed(input, accessor, frame, results, true);
return new ObjectArrayColumn(results, accessor.getType());
}
);
}

@Override
public boolean validateEquivalent(Processor otherProcessor)
{
if (otherProcessor instanceof WindowFirstProcessor) {
final WindowFirstProcessor other = (WindowFirstProcessor) otherProcessor;
return Objects.equals(frame, other.frame) && intervalValidation(other);
}
return false;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,34 +20,81 @@
package org.apache.druid.query.operator.window.value;

import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.druid.java.util.common.ISE;
import org.apache.druid.error.DruidException;
import org.apache.druid.query.operator.window.Processor;
import org.apache.druid.query.operator.window.WindowFrame;
import org.apache.druid.query.rowsandcols.RowsAndColumns;
import org.apache.druid.query.rowsandcols.column.ColumnAccessor;
import org.apache.druid.query.rowsandcols.column.ConstantObjectColumn;
import org.apache.druid.query.rowsandcols.column.ObjectArrayColumn;

import javax.annotation.Nullable;
import java.util.Objects;

public class WindowLastProcessor extends WindowValueProcessorBase
{
@Nullable
private final WindowFrame frame;

@JsonCreator
public WindowLastProcessor(
@JsonProperty("inputColumn") String inputColumn,
@JsonProperty("outputColumn") String outputColumn
@JsonProperty("outputColumn") String outputColumn,
@JsonProperty("frame") @Nullable WindowFrame frame
)
{
super(inputColumn, outputColumn);
this.frame = frame;
}

@Nullable
@JsonProperty("frame")
@JsonInclude(JsonInclude.Include.NON_NULL)
public WindowFrame getFrame()
{
return frame;
}

@Override
public RowsAndColumns process(RowsAndColumns input)
{
final int lastIndex = input.numRows() - 1;
if (lastIndex < 0) {
throw new ISE("Called with an input partition of size 0. The call site needs to not do that.");
final int numRows = input.numRows();

if (numRows == 0) {
throw DruidException.defensive("Called with an input partition of size 0. The call site needs to not do that.");
}

return processInternal(input, column -> {
final ColumnAccessor accessor = column.toAccessor();
return new ConstantObjectColumn(accessor.getObject(lastIndex), accessor.numRows(), accessor.getType());
});
if (frame == null) {
final int lastIndex = numRows - 1;
return processInternal(
input,
column -> {
final ColumnAccessor accessor = column.toAccessor();
return new ConstantObjectColumn(accessor.getObject(lastIndex), accessor.numRows(), accessor.getType());
}
);
}

return processInternal(
input,
column -> {
final ColumnAccessor accessor = column.toAccessor();
final Object[] results = new Object[numRows];
computeFirstOrLastValueFramed(input, accessor, frame, results, false);
return new ObjectArrayColumn(results, accessor.getType());
}
);
}

@Override
public boolean validateEquivalent(Processor otherProcessor)
{
if (otherProcessor instanceof WindowLastProcessor) {
final WindowLastProcessor other = (WindowLastProcessor) otherProcessor;
return Objects.equals(frame, other.frame) && intervalValidation(other);
}
return false;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,11 +20,15 @@
package org.apache.druid.query.operator.window.value;

import com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.druid.error.DruidException;
import org.apache.druid.java.util.common.ISE;
import org.apache.druid.query.operator.window.Processor;
import org.apache.druid.query.operator.window.WindowFrame;
import org.apache.druid.query.rowsandcols.RowsAndColumns;
import org.apache.druid.query.rowsandcols.column.Column;
import org.apache.druid.query.rowsandcols.column.ColumnAccessor;
import org.apache.druid.query.rowsandcols.semantic.AppendableRowsAndColumns;
import org.apache.druid.query.rowsandcols.semantic.ClusteredGroupPartitioner;

import java.util.Collections;
import java.util.List;
Expand Down Expand Up @@ -108,4 +112,63 @@ public List<String> getOutputColumnNames()
{
return Collections.singletonList(outputColumn);
}

/**
* Computes first or last value for each row respecting the given window frame. For each row, selects
* either the first or last element of the frame. Entries for rows where the frame is empty are left as null.
*
* @param first if true, selects the first element of the frame; if false, selects the last
*/
static void computeFirstOrLastValueFramed(
final RowsAndColumns rac,
final ColumnAccessor accessor,
final WindowFrame frame,
final Object[] results,
final boolean first
)
{
final int numRows = accessor.numRows();
final WindowFrame.Rows rowsFrame = frame.unwrap(WindowFrame.Rows.class);
if (rowsFrame != null) {
final int lower = rowsFrame.getLowerOffsetClamped(numRows);
final int upper = rowsFrame.getUpperOffsetClamped(numRows);
for (int i = 0; i < numRows; i++) {
final int frameLower = Math.max(0, i + lower);
final int frameUpper = Math.min(numRows - 1, i + upper);
if (frameLower <= frameUpper) {
results[i] = accessor.getObject(first ? frameLower : frameUpper);
}
}
return;
}

// Note that this logic would not be correct for SQL like:
// FIRST_VALUE(x) OVER (ORDER BY t RANGE BETWEEN 1 PRECEDING AND CURRENT ROW)
// The SQL validator will reject queries with offset-based RANGE frames such as this.
// See https://github.com/apache/druid/issues/15767 for more details.
final WindowFrame.Groups groupsFrame = frame.unwrap(WindowFrame.Groups.class);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Bounded RANGE frames still execute with GROUPS semantics

The new framed path treats every non-ROWS frame as WindowFrame.Groups, and Windowing maps Calcite's isRows == false windows into that representation. That is not correct for bounded SQL RANGE frames, which must use value-distance on the ORDER BY expression. Queries like FIRST_VALUE(x) OVER (ORDER BY t RANGE BETWEEN 1 PRECEDING AND CURRENT ROW) will still return wrong answers even though this PR now enables framed FIRST_VALUE/LAST_VALUE.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Processing won't get this far due to validation. There's a test for it in CalciteQueryTest#testUnSupportedRangeBounds, expecting an error message like:

Order By with RANGE clause currently supports only UNBOUNDED or CURRENT ROW.

I added a comment here about that.

if (groupsFrame != null) {
final int[] boundaries =
ClusteredGroupPartitioner.fromRAC(rac).computeBoundaries(groupsFrame.getOrderByColumns());
final int numGroups = boundaries.length - 1;
final int lower = groupsFrame.getLowerOffsetClamped(numGroups);
final int upper = groupsFrame.getUpperOffsetClamped(numGroups);

for (int g = 0; g < numGroups; g++) {
final int lowerGroup = Math.max(0, g + lower);
final int upperGroup = Math.min(numGroups - 1, g + upper);
if (lowerGroup <= upperGroup) {
final Object value = first
? accessor.getObject(boundaries[lowerGroup])
: accessor.getObject(boundaries[upperGroup + 1] - 1);
for (int i = boundaries[g]; i < boundaries[g + 1]; i++) {
results[i] = value;
}
}
}
return;
}

throw DruidException.defensive("Unable to handle WindowFrame[%s]", frame);
}
}
Loading
Loading