Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
604620a
working on options for excluding expired data
pri1712 Jul 18, 2026
8aa1b3d
working on options for excluding expired data
pri1712 Jul 18, 2026
9ad8fa2
working on options for excluding expired data
pri1712 Jul 18, 2026
26e0f08
working on options for excluding expired data
pri1712 Jul 18, 2026
022849d
tests for excluding expired data
pri1712-fl Jul 18, 2026
d2270e9
working on options for excluding expired data
pri1712 Jul 18, 2026
19218c4
Merge pull request #3 from pri1712/pri1712/excludeexpireddata
pri1712 Jul 18, 2026
6f06fd5
Merge branch 'apache:master' into master
pri1712 Jul 18, 2026
93361d1
Merge branch 'apache:master' into master
pri1712 Jul 19, 2026
6f0883f
adding integration tests
pri1712 Aug 2, 2026
478b20c
Merge branch 'apache:master' into master
pri1712 Aug 2, 2026
99fbdc8
Merge pull request #4 from pri1712/pri1712/excludeexpireddata
pri1712 Aug 2, 2026
670de85
Merge pull request #5 from pri1712/master
pri1712 Aug 2, 2026
97483bd
fixing various issues with checkstyle and integ tests
pri1712 Aug 2, 2026
5f036a0
fixing various issues with checkstyle and integ tests
pri1712 Aug 2, 2026
452c349
fixing various issues with checkstyle and integ tests
pri1712 Aug 2, 2026
dd6c64f
fixing checkstyle violations
pri1712 Aug 2, 2026
42cb303
Merge branch 'master' of https://github.com/pri1712/pinot
pri1712-fl Sep 12, 2026
81cb1ab
Merge branch 'pri1712/excludeexpireddata' of https://github.com/pri17…
pri1712-fl Sep 12, 2026
6973555
query opts for skipping expired records
pri1712 Sep 19, 2026
551556c
cleaning up
pri1712 Sep 19, 2026
1319329
Merge branch 'pri1712/excludeexpireddata' of https://github.com/pri17…
pri1712-fl Sep 19, 2026
c3734e6
skipping expired records
pri1712-fl Sep 19, 2026
eaf88cc
fixing checkstyle issues
pri1712 Sep 19, 2026
118126b
fixing tests and nitpicks
pri1712 Sep 19, 2026
7a54a68
Update log4j2-test.xml
pri1712 Sep 19, 2026
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 @@ -116,7 +116,11 @@
import org.apache.pinot.spi.config.table.FieldConfig;
import org.apache.pinot.spi.config.table.QueryConfig;
import org.apache.pinot.spi.config.table.RoutingConfig;
import org.apache.pinot.spi.config.table.SegmentsValidationAndRetentionConfig;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.data.DateTimeFieldSpec;
import org.apache.pinot.spi.data.DateTimeFieldSpec.TimeFormat;
import org.apache.pinot.spi.data.DateTimeFormatSpec;
import org.apache.pinot.spi.data.LogicalTableConfig;
import org.apache.pinot.spi.data.Schema;
import org.apache.pinot.spi.env.PinotConfiguration;
Expand Down Expand Up @@ -708,12 +712,16 @@ protected BrokerResponse doHandleRequest(long requestId, String query, SqlNodeAn
BrokerRequest offlineBrokerRequest = null;
BrokerRequest realtimeBrokerRequest = null;

boolean skipExpiredRecords = QueryOptionsUtils.isSkipExpiredRecords(serverPinotQuery.getQueryOptions());
if (routeInfo.isHybrid()) {
// Hybrid
PinotQuery offlinePinotQuery = serverPinotQuery.deepCopy();
offlinePinotQuery.getDataSource().setTableName(offlineTableName);
assert timeBoundaryInfo != null;
attachTimeBoundary(offlinePinotQuery, timeBoundaryInfo, true);
if (skipExpiredRecords) {
handleSkipExpiredRecords(offlineTableConfig, schema, offlinePinotQuery);
}
handleExpressionOverride(offlinePinotQuery, _tableCache.getExpressionOverrideMap(offlineTableName));
handleTimestampIndexOverride(offlinePinotQuery, offlineTableConfig);
// Re-optimize after attaching the time boundary filter so that filter optimizers (e.g. NumericalFilterOptimizer,
Expand All @@ -724,6 +732,9 @@ protected BrokerResponse doHandleRequest(long requestId, String query, SqlNodeAn
PinotQuery realtimePinotQuery = serverPinotQuery.deepCopy();
realtimePinotQuery.getDataSource().setTableName(realtimeTableName);
attachTimeBoundary(realtimePinotQuery, timeBoundaryInfo, false);
if (skipExpiredRecords) {
handleSkipExpiredRecords(realtimeTableConfig, schema, realtimePinotQuery);
}
handleExpressionOverride(realtimePinotQuery, _tableCache.getExpressionOverrideMap(realtimeTableName));
handleTimestampIndexOverride(realtimePinotQuery, realtimeTableConfig);
_queryOptimizer.optimize(realtimePinotQuery, schema);
Expand All @@ -735,6 +746,10 @@ protected BrokerResponse doHandleRequest(long requestId, String query, SqlNodeAn
} else if (routeInfo.isOffline()) {
// OFFLINE only
setTableName(serverBrokerRequest, offlineTableName);
if (skipExpiredRecords) {
handleSkipExpiredRecords(offlineTableConfig, schema, serverPinotQuery);
_queryOptimizer.optimize(serverPinotQuery, schema);

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.

one query i have is, is it okay to reoptimize here? @xiangfu0 ?

for a per table type retention field, this seems like the best place to put it

}
handleExpressionOverride(serverPinotQuery, _tableCache.getExpressionOverrideMap(offlineTableName));
handleTimestampIndexOverride(serverPinotQuery, offlineTableConfig);
offlineBrokerRequest = serverBrokerRequest;
Expand All @@ -744,6 +759,10 @@ protected BrokerResponse doHandleRequest(long requestId, String query, SqlNodeAn
} else {
// REALTIME only
setTableName(serverBrokerRequest, realtimeTableName);
if (skipExpiredRecords) {
handleSkipExpiredRecords(realtimeTableConfig, schema, serverPinotQuery);
_queryOptimizer.optimize(serverPinotQuery, schema);
}
handleExpressionOverride(serverPinotQuery, _tableCache.getExpressionOverrideMap(realtimeTableName));
handleTimestampIndexOverride(serverPinotQuery, realtimeTableConfig);
realtimeBrokerRequest = serverBrokerRequest;
Expand Down Expand Up @@ -1246,7 +1265,6 @@ private CompileResult compileRequest(long requestId, String query, SqlNodeAndOpt
if (_enableDistinctCountBitmapOverride) {
handleDistinctCountBitmapOverride(serverPinotQuery);
}

Schema schema = _tableCache.getSchema(rawTableName);
_queryOptimizer.optimize(serverPinotQuery, schema);

Expand Down Expand Up @@ -1861,6 +1879,78 @@ private static void handleDistinctCountBitmapOverride(Expression expression) {
}
}

/// Attaches a `timeColumn >= (now - retention)` filter to the given query so records outside the table's retention
/// window are excluded, even if their segment has not yet been deleted (see issue #16689). Applied per-leg for hybrid
/// tables so the offline and realtime sides each use their own retention. No-ops (with a debug log) when the config,
/// time column, retention, or schema spec is missing/malformed, rather than failing the query.
@VisibleForTesting
static void handleSkipExpiredRecords(@Nullable TableConfig tableConfig, @Nullable Schema schema,
PinotQuery pinotQuery) {
if (tableConfig == null || schema == null) {
return;
}
String tableNameWithType = tableConfig.getTableName();
SegmentsValidationAndRetentionConfig validationConfig = tableConfig.getValidationConfig();
if (validationConfig == null) {
LOGGER.debug("skipExpiredRecords: no validation config for table {}, skipping retention filter",
tableNameWithType);
return;
}

String timeColumnName = validationConfig.getTimeColumnName();
if (timeColumnName == null) {
LOGGER.debug("skipExpiredRecords: no time column configured for table {}, skipping retention filter",
tableNameWithType);
return;
}

Long retentionMs = getRetentionMs(validationConfig);
if (retentionMs == null) {
LOGGER.debug("skipExpiredRecords: no valid retention configured for table {}, skipping retention filter",
tableNameWithType);
return;
}
long cutOffMs = System.currentTimeMillis() - retentionMs;

DateTimeFieldSpec timeFieldSpec = schema.getSpecForTimeColumn(timeColumnName);
if (timeFieldSpec == null) {
LOGGER.debug("skipExpiredRecords: time column {} not found in schema for table {}, skipping retention filter",
timeColumnName, tableNameWithType);
return;
}

DateTimeFormatSpec formatSpec = timeFieldSpec.getFormatSpec();
String cutOffValue = formatSpec.fromMillisToFormat(cutOffMs);
Expression cutOffLiteral = formatSpec.getTimeFormat() == TimeFormat.EPOCH
? RequestUtils.getLiteralExpression(Long.parseLong(cutOffValue))
: RequestUtils.getLiteralExpression(cutOffValue);
Expression retentionFilter = RequestUtils.getFunctionExpression(FilterKind.GREATER_THAN_OR_EQUAL.name(),
RequestUtils.getIdentifierExpression(timeColumnName), cutOffLiteral);

Expression existingFilter = pinotQuery.getFilterExpression();
pinotQuery.setFilterExpression(existingFilter != null
? RequestUtils.getFunctionExpression(FilterKind.AND.name(), existingFilter, retentionFilter)
: retentionFilter);
LOGGER.debug("skipExpiredRecords: attached retention filter {} >= {} (cutOffMs={}) for table {}", timeColumnName,
cutOffValue, cutOffMs, tableNameWithType);
}

/// Parses the retention window in millis from the validation config, or `null` when retention is not configured or is
/// malformed (in which case no retention filter is applied rather than failing the query).
@Nullable
private static Long getRetentionMs(SegmentsValidationAndRetentionConfig validationConfig) {
String retentionUnit = validationConfig.getRetentionTimeUnit();
String retentionValue = validationConfig.getRetentionTimeValue();
if (StringUtils.isEmpty(retentionUnit) || StringUtils.isEmpty(retentionValue)) {
return null;
}
try {
return TimeUnit.valueOf(retentionUnit.toUpperCase()).toMillis(Long.parseLong(retentionValue));
} catch (IllegalArgumentException e) {
return null;
}
}

private HandlerContext getHandlerContext(@Nullable QueryConfig offlineTableQueryConfig,
@Nullable QueryConfig realtimeTableQueryConfig) {
Boolean disableGroovyOverride = null;
Expand Down Expand Up @@ -2630,6 +2720,8 @@ private ImplicitHybridTableRouteInfo prepareBaseTableHybridRoute(BrokerRequest b
TableConfig realtimeTableConfig = baseRouteInfo.getRealtimeTableConfig();
TimeBoundaryInfo timeBoundaryInfo = baseRouteInfo.getTimeBoundaryInfo();

boolean skipExpiredRecords =
QueryOptionsUtils.isSkipExpiredRecords(baseBrokerRequest.getPinotQuery().getQueryOptions());
if (baseRouteInfo.isHybrid()) {
PinotQuery basePinotQuery = baseBrokerRequest.getPinotQuery();

Expand All @@ -2638,6 +2730,9 @@ private ImplicitHybridTableRouteInfo prepareBaseTableHybridRoute(BrokerRequest b
if (timeBoundaryInfo != null) {
attachTimeBoundary(offlinePinotQuery, timeBoundaryInfo, true);
}
if (skipExpiredRecords) {
handleSkipExpiredRecords(offlineTableConfig, schema, offlinePinotQuery);
}
handleExpressionOverride(offlinePinotQuery, _tableCache.getExpressionOverrideMap(offlineTableName));
handleTimestampIndexOverride(offlinePinotQuery, offlineTableConfig);
_queryOptimizer.optimize(offlinePinotQuery, schema);
Expand All @@ -2648,6 +2743,9 @@ private ImplicitHybridTableRouteInfo prepareBaseTableHybridRoute(BrokerRequest b
if (timeBoundaryInfo != null) {
attachTimeBoundary(realtimePinotQuery, timeBoundaryInfo, false);
}
if (skipExpiredRecords) {
handleSkipExpiredRecords(realtimeTableConfig, schema, realtimePinotQuery);
}
handleExpressionOverride(realtimePinotQuery, _tableCache.getExpressionOverrideMap(realtimeTableName));
handleTimestampIndexOverride(realtimePinotQuery, realtimeTableConfig);
_queryOptimizer.optimize(realtimePinotQuery, schema);
Expand All @@ -2657,12 +2755,20 @@ private ImplicitHybridTableRouteInfo prepareBaseTableHybridRoute(BrokerRequest b
hybridRoute.setRealtimeBrokerRequest(realtimeBrokerRequest);
} else if (baseRouteInfo.isOffline()) {
setTableName(baseBrokerRequest, offlineTableName);
if (skipExpiredRecords) {
handleSkipExpiredRecords(offlineTableConfig, schema, baseBrokerRequest.getPinotQuery());
_queryOptimizer.optimize(baseBrokerRequest.getPinotQuery(), schema);
}
handleExpressionOverride(baseBrokerRequest.getPinotQuery(),
_tableCache.getExpressionOverrideMap(offlineTableName));
handleTimestampIndexOverride(baseBrokerRequest.getPinotQuery(), offlineTableConfig);
hybridRoute.setOfflineBrokerRequest(baseBrokerRequest);
} else {
setTableName(baseBrokerRequest, realtimeTableName);
if (skipExpiredRecords) {
handleSkipExpiredRecords(realtimeTableConfig, schema, baseBrokerRequest.getPinotQuery());
_queryOptimizer.optimize(baseBrokerRequest.getPinotQuery(), schema);
}
handleExpressionOverride(baseBrokerRequest.getPinotQuery(),
_tableCache.getExpressionOverrideMap(realtimeTableName));
handleTimestampIndexOverride(baseBrokerRequest.getPinotQuery(), realtimeTableConfig);
Expand Down
Loading
Loading