Jackie-Jiang commented on code in PR #19011:
URL: https://github.com/apache/pinot/pull/19011#discussion_r4113625619
##########
pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java:
##########
@@ -708,12 +712,16 @@ protected BrokerResponse doHandleRequest(long requestId,
String query, SqlNodeAn
BrokerRequest offlineBrokerRequest = null;
BrokerRequest realtimeBrokerRequest = null;
+ boolean skipExpiredRecords =
QueryOptionsUtils.isSkipExpiredRecords(serverPinotQuery.getQueryOptions());
Review Comment:
**Retention is not preserved through materialized-view rewrites.** A
FULL_REWRITE can replace `serverPinotQuery` with the MV query before this
option is checked, so the later cutoff comes from the MV table config rather
than the queried base table. In SPLIT_REWRITE, `prepareBaseTableHybridRoute`
adds the cutoff to base legs, while the MV leg is dispatched with only its
watermark bound. With five-day base retention and an MV holding 30-day-old
rows, `skipExpiredRecords=true` can still return expired data. Please carry the
source-table cutoff through both rewrite modes, or suppress MV rewrites for
this option until they preserve its semantics.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java:
##########
@@ -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);
Review Comment:
**Formatted STRING dates need chronological comparison.** The earlier
string-ordering concern remains in this implementation: `SIMPLE_DATE_FORMAT`
allows patterns such as `MM/dd/yyyy`, but the server evaluates a STRING `>=`
with lexical `String.compareTo`. For cutoff `09/26/2026`, an expired value
`12/01/2025` passes this filter; other patterns can exclude fresh rows. Please
compare normalized time values or restrict the option to formats proven to
preserve chronological ordering, with a regression case.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java:
##########
@@ -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);
Review Comment:
**Apply retention per physical table on logical routes.**
`LogicalTableRouteProvider` gets this config from `refOfflineTableName` and
sends the same filtered request to every physical offline table. Logical-table
validation does not require those tables to have equal retention. If the
reference retains 30 days and another table retains seven, rows aged 7–30 days
from the latter remain visible; reversing the values excludes valid rows.
Please derive the cutoff for each physical scan or define and validate a shared
logical-table retention policy.
##########
pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java:
##########
@@ -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);
Review Comment:
**Round the lower bound up for coarse epoch formats.**
`fromMillisToFormat()` truncates to the time column’s unit and size. For
`1:HOURS:EPOCH`, a cutoff at 12:30 produces the literal for 12:00, so `>=`
includes a record timestamped 12:00 even though it is expired. Use the first
representable column value at or after the cutoff when the cutoff falls inside
a bucket, and cover this case in a test.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]