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

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -307,13 +307,13 @@ public void testPushDownFilterIntoWindow() {
String[] expectedHeader = new String[] {"time", "device", "value", "rn"};
String[] retArray =
new String[] {
"2021-01-01T09:10:00.000Z,d1,1.0,1,",
"2021-01-01T09:05:00.000Z,d1,3.0,2,",
"2021-01-01T09:10:00.000Z,d1,1.0,1,",
"2021-01-01T09:08:00.000Z,d2,2.0,1,",
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT * FROM (SELECT *, row_number() OVER (PARTITION BY device ORDER BY value) as rn FROM demo) WHERE rn <= 2 ORDER BY device, time",
"SELECT * FROM (SELECT *, row_number() OVER (PARTITION BY device ORDER BY value, time) as rn FROM demo) WHERE rn <= 2 ORDER BY device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -327,7 +327,7 @@ public void testPushDownLimitIntoWindow() {
"2021-01-01T09:05:00.000Z,d1,3.0,2,", "2021-01-01T09:07:00.000Z,d1,5.0,4,",
};
tableResultSetEqualTest(
"SELECT * FROM (SELECT *, row_number() OVER (PARTITION BY device ORDER BY value) as rn FROM demo) ORDER BY device, time LIMIT 2 ",
"SELECT * FROM (SELECT *, row_number() OVER (PARTITION BY device ORDER BY value, time) as rn FROM demo) ORDER BY device, time LIMIT 2 ",
expectedHeader,
retArray,
DATABASE_NAME);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,7 @@ public void testOrderBy() {
"2021-01-01T09:07:00.000Z,d1,5.0,6,",
};
tableResultSetEqualTest(
"SELECT *, count(value) OVER (ORDER BY value) AS cnt FROM demo",
"SELECT *, count(value) OVER (ORDER BY value) AS cnt FROM demo ORDER BY value, device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -216,7 +216,7 @@ public void testOrderByWithNulls() {
"2021-01-01T09:20:00.000Z,null,null,8,",
};
tableResultSetEqualTest(
"SELECT *, count(value) OVER (ORDER BY value) AS cnt FROM demo2 ORDER BY value, device",
"SELECT *, count(value) OVER (ORDER BY value) AS cnt FROM demo2 ORDER BY value, device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -235,7 +235,7 @@ public void testPartitionByAndOrderBy() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, rank() OVER (PARTITION BY device ORDER BY value) AS rnk FROM demo ORDER BY device",
"SELECT *, rank() OVER (PARTITION BY device ORDER BY value) AS rnk FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand Down Expand Up @@ -278,7 +278,7 @@ public void testRowsFraming() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY time ROWS 1 PRECEDING) AS cnt FROM demo ORDER BY device",
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY time ROWS 1 PRECEDING) AS cnt FROM demo ORDER BY device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -297,7 +297,7 @@ public void testGroupsFraming() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY value GROUPS BETWEEN 1 PRECEDING AND CURRENT ROW) AS cnt FROM demo ORDER BY device",
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY value GROUPS BETWEEN 1 PRECEDING AND CURRENT ROW) AS cnt FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -316,7 +316,7 @@ public void testRangeFraming() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY value RANGE BETWEEN 2 PRECEDING AND CURRENT ROW) AS cnt FROM demo ORDER BY device",
"SELECT *, count(value) OVER (PARTITION BY device ORDER BY value RANGE BETWEEN 2 PRECEDING AND CURRENT ROW) AS cnt FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -335,7 +335,7 @@ public void testAggregation() {
"2021-01-01T09:15:00.000Z,d2,4.0,6.0,",
};
tableResultSetEqualTest(
"SELECT *, sum(value) OVER (PARTITION BY device ORDER BY value) AS sum FROM demo ORDER BY device",
"SELECT *, sum(value) OVER (PARTITION BY device ORDER BY value) AS sum FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -354,7 +354,7 @@ public void testFirstValue() {
"2021-01-01T09:15:00.000Z,d2,4.0,2.0,",
};
tableResultSetEqualTest(
"SELECT *, first_value(value) OVER (PARTITION BY device ORDER BY value ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS fv FROM demo ORDER BY device",
"SELECT *, first_value(value) OVER (PARTITION BY device ORDER BY value, time ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS fv FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -373,7 +373,7 @@ public void testLastValue() {
"2021-01-01T09:15:00.000Z,d2,4.0,4.0,",
};
tableResultSetEqualTest(
"SELECT *, last_value(value) OVER (PARTITION BY device ORDER BY value ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS lv FROM demo ORDER BY device",
"SELECT *, last_value(value) OVER (PARTITION BY device ORDER BY value, time ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS lv FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -392,7 +392,7 @@ public void testNthValue() {
"2021-01-01T09:15:00.000Z,d2,4.0,4.0,",
};
tableResultSetEqualTest(
"SELECT *, nth_value(value, 2) OVER (PARTITION BY device ORDER BY value ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS nv FROM demo ORDER BY device",
"SELECT *, nth_value(value, 2) OVER (PARTITION BY device ORDER BY value, time ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING) AS nv FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -411,7 +411,7 @@ public void testLead() {
"2021-01-01T09:15:00.000Z,d2,4.0,null,",
};
tableResultSetEqualTest(
"SELECT *, lead(value) OVER (PARTITION BY device ORDER BY time) AS ld FROM demo ORDER BY device",
"SELECT *, lead(value) OVER (PARTITION BY device ORDER BY time) AS ld FROM demo ORDER BY device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -435,7 +435,7 @@ public void testLag() {
"2021-01-01T09:15:00.000Z,d2,4.0,2.0,",
};
tableResultSetEqualTest(
"SELECT *, lag(value) OVER (PARTITION BY device ORDER BY time) AS lg FROM demo ORDER BY device",
"SELECT *, lag(value) OVER (PARTITION BY device ORDER BY time) AS lg FROM demo ORDER BY device, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -459,7 +459,7 @@ public void testRank() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, rank() OVER (PARTITION BY device ORDER BY value) AS rk FROM demo ORDER BY device",
"SELECT *, rank() OVER (PARTITION BY device ORDER BY value) AS rk FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -478,7 +478,7 @@ public void testDenseRank() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, dense_rank() OVER (PARTITION BY device ORDER BY value) AS rk FROM demo ORDER BY device",
"SELECT *, dense_rank() OVER (PARTITION BY device ORDER BY value) AS rk FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -497,7 +497,7 @@ public void testRowNumber() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, row_number() OVER (PARTITION BY device ORDER BY value) AS rn FROM demo ORDER BY device",
"SELECT *, row_number() OVER (PARTITION BY device ORDER BY value, time) AS rn FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -516,7 +516,7 @@ public void testPercentRank() {
"2021-01-01T09:15:00.000Z,d2,4.0,1.0,",
};
tableResultSetEqualTest(
"SELECT *, percent_rank() OVER (PARTITION BY device ORDER BY value) AS pr FROM demo ORDER BY device",
"SELECT *, percent_rank() OVER (PARTITION BY device ORDER BY value) AS pr FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -535,7 +535,7 @@ public void testCumeDist() {
"2021-01-01T09:15:00.000Z,d2,4.0,1.0,",
};
tableResultSetEqualTest(
"SELECT *, cume_dist() OVER (PARTITION BY device ORDER BY value) AS cd FROM demo ORDER BY device",
"SELECT *, cume_dist() OVER (PARTITION BY device ORDER BY value) AS cd FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand All @@ -554,7 +554,7 @@ public void testNTile() {
"2021-01-01T09:15:00.000Z,d2,4.0,2,",
};
tableResultSetEqualTest(
"SELECT *, ntile(2) OVER (PARTITION BY device ORDER BY value) AS nt FROM demo ORDER BY device",
"SELECT *, ntile(2) OVER (PARTITION BY device ORDER BY value, time) AS nt FROM demo ORDER BY device, value, time",
expectedHeader,
retArray,
DATABASE_NAME);
Expand Down Expand Up @@ -635,7 +635,7 @@ public void testComplexQuery() {
+ " DATA => (SELECT time, flow FROM demo3 WHERE device = 'd0'),\n"
+ " COL => 'flow',\n"
+ " DELTA => 0.0)\n"
+ " GROUP BY flow",
+ " GROUP BY flow ORDER BY flow",
expectedHeader,
retArray,
DATABASE_NAME);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,9 @@
import java.sql.SQLException;
import java.sql.Statement;
import java.sql.Types;
import java.time.Instant;
import java.util.Arrays;
import java.util.Comparator;
import java.util.List;

import static org.apache.iotdb.db.it.utils.TestUtils.*;
Expand Down Expand Up @@ -144,6 +146,8 @@ public void testHopFunction() {
expectedHeader,
retArray,
DATABASE_NAME);
assertTimeWindowOrdering(
"HOP(DATA => bid, TIMECOL => 'time', SLIDE => 5m, SIZE => 10m)", expectedHeader, retArray);

expectedHeader = new String[] {"window_start", "window_end", "stock_id", "sum"};
retArray =
Expand Down Expand Up @@ -503,6 +507,22 @@ public void testCapacityFunction() {
expectedHeader,
retArray,
DATABASE_NAME);
assertSortedRows(
"SELECT * FROM CAPACITY(DATA => bid PARTITION BY stock_id ORDER BY time, SIZE => 2, SLIDE => 1) ORDER BY stock_id DESC, time DESC, window_index DESC",
expectedHeader,
retArray,
Comparator.comparing((String[] row) -> row[2])
.thenComparing(row -> Instant.parse(row[1]))
.thenComparingLong(row -> Long.parseLong(row[0]))
.reversed());
tableResultSetEqualTest(
"SELECT stock_id FROM CAPACITY(DATA => bid PARTITION BY stock_id ORDER BY time, SIZE => 2, SLIDE => 1) ORDER BY stock_id",
new String[] {"stock_id"},
new String[] {
"AAPL,", "AAPL,", "AAPL,", "AAPL,", "AAPL,",
"TESL,", "TESL,", "TESL,", "TESL,", "TESL,"
},
DATABASE_NAME);

// CAPACITY with SIZE=3, SLIDE=2 (overlapping windows, different params)
expectedHeader = new String[] {"window_index", "time", "stock_id", "price", "s1"};
Expand Down Expand Up @@ -587,6 +607,8 @@ public void testTumbleFunction() {
expectedHeader,
retArray,
DATABASE_NAME);
assertTimeWindowOrdering(
"TUMBLE(DATA => bid, TIMECOL => 'time', SIZE => 10m)", expectedHeader, retArray);

// TUMBLE (10m) + origin
expectedHeader = new String[] {"window_start", "window_end", "time", "stock_id", "price", "s1"};
Expand Down Expand Up @@ -653,6 +675,10 @@ public void testCumulateFunction() {
expectedHeader,
retArray,
DATABASE_NAME);
assertTimeWindowOrdering(
"CUMULATE(DATA => bid, TIMECOL => 'time', STEP => 6m, SIZE => 12m)",
expectedHeader,
retArray);

expectedHeader = new String[] {"window_start", "window_end", "time", "stock_id", "price", "s1"};
retArray =
Expand Down Expand Up @@ -1595,6 +1621,57 @@ public void testXCorrRejectsUnexpectedCalculationColumnCount() {
DATABASE_NAME);
}

private static void assertTimeWindowOrdering(
String function, String[] expectedHeader, String[] expectedRows) {
// A row can belong to several windows. Device/time no longer uniquely identifies a TVF row,
// so the descending window keys must still be sorted, even after all TAGs and TIME.
assertSortedRows(
"SELECT * FROM "
+ function
+ " ORDER BY stock_id, time, window_start DESC, window_end DESC",
expectedHeader,
expectedRows,
Comparator.comparing((String[] row) -> row[3])
.thenComparing(row -> Instant.parse(row[2]))
.thenComparing(row -> Instant.parse(row[0]), Comparator.reverseOrder())
.thenComparing(row -> Instant.parse(row[1]), Comparator.reverseOrder()));
assertSortedRows(
"SELECT * FROM "
+ function
+ " ORDER BY window_start DESC, stock_id, time DESC, window_end DESC",
expectedHeader,
expectedRows,
Comparator.comparing((String[] row) -> Instant.parse(row[0]), Comparator.reverseOrder())
.thenComparing(row -> row[3])
.thenComparing(row -> Instant.parse(row[2]), Comparator.reverseOrder())
.thenComparing(row -> Instant.parse(row[1]), Comparator.reverseOrder()));

// Project away generated columns to exercise elimination for a descending pass-through order.
String[] passThroughRows =
Arrays.stream(expectedRows)
.map(row -> row.split(","))
.map(row -> row[3] + "," + row[2] + "," + row[4] + "," + row[5] + ",")
.toArray(String[]::new);
assertSortedRows(
"SELECT stock_id, time, price, s1 FROM " + function + " ORDER BY stock_id DESC, time DESC",
new String[] {"stock_id", "time", "price", "s1"},
passThroughRows,
Comparator.comparing((String[] row) -> row[0])
.thenComparing(row -> Instant.parse(row[1]))
.reversed());
}

private static void assertSortedRows(
String sql, String[] header, String[] rows, Comparator<String[]> comparator) {
String[] sortedRows =
Arrays.stream(rows)
.map(row -> row.split(","))
.sorted(comparator)
.map(row -> String.join(",", row) + ",")
.toArray(String[]::new);
tableResultSetEqualTest(sql, header, sortedRows, DATABASE_NAME);
}

private static void tableResultSetEqualWithTolerance(
String sql, String[] expectedHeader, String[] expectedRetArray, String database) {
try (Connection connection = EnvFactory.getEnv().getTableConnection();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.iotdb.calc.execution.operator.source.relational.aggregation.grouped.array.LongBigArrayFIFOQueue;
import org.apache.iotdb.calc.i18n.CalcMessages;

import org.apache.tsfile.block.column.Column;
import org.apache.tsfile.read.common.block.TsBlock;
import org.apache.tsfile.utils.RamUsageEstimator;

Expand Down Expand Up @@ -307,8 +308,15 @@ public void compact() {
rowIdBuffer.setPosition(newRowIds[i], i);
}

// Compact TsBlock
// page = page.copyPositions(positionsToKeep, 0, positionsToKeep.length);
// The stored positions above now refer to the compacted block, not the original rows.
// Copy only live positions so stable row IDs keep referencing the same values.
Column[] columns = new Column[tsBlock.getValueColumnCount()];
for (int i = 0; i < columns.length; i++) {
columns[i] = tsBlock.getColumn(i).copyPositions(positionsToKeep, 0, activePositions);
}
tsBlock =
new TsBlock(
tsBlock.getTimeColumn().copyPositions(positionsToKeep, 0, activePositions), columns);
rowIds = newRowIds;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ public TsBlock next() throws Exception {

TsBlock tempResult = null;
while (tempResult == null && !cachedTsBlock.isEmpty()) {
prepareFill();
TsBlock originTsBlock = cachedTsBlock.get(0);
long currentEndRowIndex =
cachedRowIndex.get(0) + cachedLastRowIndexForNonNullHelperColumn.get(0);
Expand Down Expand Up @@ -175,6 +176,14 @@ void resetFill() {
// do nothing
}

void prepareFill() {
// Grouped fills may reset their state before preparing lookahead for a new group.
}

boolean isGroupEnd(int cachedBlockIndex) {
return false;
}

@Override
public boolean hasNext() throws Exception {
// if child.hasNext() return false, it means that there is no more tsBlocks
Expand Down Expand Up @@ -233,6 +242,10 @@ public long ramBytesUsed() {
private boolean isCachedTsBlockEnough(int columnIndex, long currentEndRowIndex) {
// next TsBlock has already been in the cachedTsBlock
while (nextTsBlockIndex[columnIndex] < cachedTsBlock.size()) {
if (isGroupEnd(nextTsBlockIndex[columnIndex] - 1)) {
// No later value belongs to this group, even if another group's blocks are cached.
return true;
}
TsBlock nextTsBlock = cachedTsBlock.get(nextTsBlockIndex[columnIndex]);
long startRowIndex = cachedRowIndex.get(nextTsBlockIndex[columnIndex]);
nextTsBlockIndex[columnIndex]++;
Expand All @@ -244,7 +257,7 @@ private boolean isCachedTsBlockEnough(int columnIndex, long currentEndRowIndex)
return true;
}
}
return false;
return isGroupEnd(nextTsBlockIndex[columnIndex] - 1);
}

/**
Expand Down
Loading
Loading