This is an automated email from the ASF dual-hosted git repository.
xiangfu0 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 71475e49ef3 add broker config options for sql log redaction (#18604)
71475e49ef3 is described below
commit 71475e49ef39e8ff40347c939b4e6caaafa9458e
Author: Johan Adami <[email protected]>
AuthorDate: Mon Jul 27 13:00:09 2026 -0400
add broker config options for sql log redaction (#18604)
* add broker config options for sql log redaction
* do not parse twice
* fix checkstyle
* default to FULL redaction when failing to parse
* redact other logs too
* remove DeepCopyShuttle
* undo unneeded changes
* add docstrings
* capture a few more unredacted logs
---------
Co-authored-by: Johan Adami <[email protected]>
---
.../apache/pinot/broker/querylog/QueryLogger.java | 65 +++++++-
.../BaseSingleStageBrokerRequestHandler.java | 93 ++++++++----
.../MultiStageBrokerRequestHandler.java | 44 ++++--
.../pinot/broker/querylog/QueryLoggerTest.java | 167 ++++++++++++++++++---
.../apache/pinot/spi/utils/CommonConstants.java | 3 +
5 files changed, 299 insertions(+), 73 deletions(-)
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/querylog/QueryLogger.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/querylog/QueryLogger.java
index 3a704c14765..c01e954b4af 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/querylog/QueryLogger.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/querylog/QueryLogger.java
@@ -48,10 +48,38 @@ public class QueryLogger {
private static final Logger LOGGER =
LoggerFactory.getLogger(QueryLogger.class);
private static final QueryLogEntry[] QUERY_LOG_ENTRY_VALUES =
QueryLogEntry.values();
+ private static final String FINGERPRINT_FAILED_QUERY_REDACTED =
"FINGERPRINT_FAILED_QUERY_REDACTED";
+ private static final String FULLY_REDACTED = "REDACTED";
+
+ public enum SqlRedactionMode {
+ // Log the full SQL query text as-is.
+ // e.g. "SELECT name FROM users WHERE id = 42 AND status = 'active'"
+ NONE,
+ // Replace literal values with placeholders using the query fingerprint,
preserving query structure.
+ // Requires query fingerprinting to be enabled (will be auto-enabled if
not configured).
+ // e.g. "SELECT name FROM users WHERE id = ? AND status = ?"
+ LITERAL_VALUES,
+ // Omit the SQL text entirely from logs, replacing it with "[REDACTED]".
+ // Use when no part of the query should appear in logs.
+ FULL;
+
+ public static SqlRedactionMode fromString(String value) {
+ try {
+ return valueOf(value.toUpperCase());
+ } catch (IllegalArgumentException e) {
+ // The default config value is NONE. If the user intended to enable
redaction but made a typo,
+ // it's safer to default to FULL instead of NONE to avoid accidentally
logging sensitive information.
+ LOGGER.warn("Invalid SQL redaction mode '{}', defaulting to FULL",
value);
+ return FULL;
+ }
+ }
+ }
+
private final int _maxQueryLengthToLog;
private final RateLimiter _logRateLimiter;
private final boolean _enableIpLogging;
private final boolean _logBeforeProcessing;
+ private final SqlRedactionMode _sqlRedactionMode;
private final Logger _logger;
private final RateLimiter _droppedLogRateLimiter;
private final AtomicLong _numDroppedLogs = new AtomicLong(0L);
@@ -63,20 +91,24 @@ public class QueryLogger {
config.getProperty(Broker.CONFIG_OF_BROKER_REQUEST_CLIENT_IP_LOGGING,
Broker.DEFAULT_BROKER_REQUEST_CLIENT_IP_LOGGING),
config.getProperty(Broker.CONFIG_OF_BROKER_QUERY_LOG_BEFORE_PROCESSING,
- Broker.DEFAULT_BROKER_QUERY_LOG_BEFORE_PROCESSING), LOGGER,
RateLimiter.create(1)
+ Broker.DEFAULT_BROKER_QUERY_LOG_BEFORE_PROCESSING),
+
SqlRedactionMode.fromString(config.getProperty(Broker.CONFIG_OF_BROKER_QUERY_LOG_SQL_REDACTION,
+ Broker.DEFAULT_BROKER_QUERY_LOG_SQL_REDACTION)),
+ LOGGER, RateLimiter.create(1)
// log once a second for dropped log count
);
}
@VisibleForTesting
QueryLogger(RateLimiter logRateLimiter, int maxQueryLengthToLog, boolean
enableIpLogging, boolean logBeforeProcessing,
- Logger logger, RateLimiter droppedLogRateLimiter) {
+ SqlRedactionMode sqlRedactionMode, Logger logger, RateLimiter
droppedLogRateLimiter) {
_logRateLimiter = logRateLimiter;
_maxQueryLengthToLog = maxQueryLengthToLog;
_enableIpLogging = enableIpLogging;
_logger = logger;
_droppedLogRateLimiter = droppedLogRateLimiter;
_logBeforeProcessing = logBeforeProcessing;
+ _sqlRedactionMode = sqlRedactionMode;
}
/**
@@ -86,15 +118,16 @@ public class QueryLogger {
*
* @param requestId the request ID
* @param query the SQL query
+ * @param queryFingerprint the query fingerprint (used when redaction is
enabled)
* @return true if the rate limiter allowed this query (not rate-limited),
false if rate-limited
*/
- public boolean logQueryReceived(long requestId, String query) {
+ public boolean logQueryReceived(long requestId, String query, @Nullable
QueryFingerprint queryFingerprint) {
if (!checkRateLimiter()) {
return false;
}
if (_logBeforeProcessing) {
- _logger.info("SQL query for request {}: {}", requestId, query);
+ _logger.info("SQL query for request {}: {}", requestId,
redactQuery(query, queryFingerprint));
}
tryLogDropped();
@@ -123,8 +156,8 @@ public class QueryLogger {
}
// always log the query last - don't add this to the QueryLogEntry enum
- queryLogBuilder.append("query=")
- .append(StringUtils.substring(params._requestContext.getQuery(), 0,
_maxQueryLengthToLog));
+ String redacted = redactQuery(params._requestContext.getQuery(),
params._requestContext.getQueryFingerprint());
+ queryLogBuilder.append("query=").append(StringUtils.substring(redacted, 0,
_maxQueryLengthToLog));
_logger.info(queryLogBuilder.toString());
tryLogDropped();
@@ -158,6 +191,26 @@ public class QueryLogger {
return _logRateLimiter.getRate();
}
+ public SqlRedactionMode getSqlRedactionMode() {
+ return _sqlRedactionMode;
+ }
+
+ public String redactQuery(String query) {
+ return redactQuery(query, null);
+ }
+
+ public String redactQuery(String query, @Nullable QueryFingerprint
queryFingerprint) {
+ switch (_sqlRedactionMode) {
+ case FULL:
+ return FULLY_REDACTED;
+ case LITERAL_VALUES:
+ return queryFingerprint != null ? queryFingerprint.getFingerprint() :
FINGERPRINT_FAILED_QUERY_REDACTED;
+ case NONE:
+ default:
+ return query;
+ }
+ }
+
private boolean shouldForceLog(@Nullable QueryLogParams params) {
if (params == null) {
return false;
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java
index af4e6b8e6c9..d1c9bc04180 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/BaseSingleStageBrokerRequestHandler.java
@@ -210,8 +210,15 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
_enableMultistageMigrationMetric =
_config.getProperty(Broker.CONFIG_OF_BROKER_ENABLE_MULTISTAGE_MIGRATION_METRIC,
Broker.DEFAULT_ENABLE_MULTISTAGE_MIGRATION_METRIC);
- _enableQueryFingerprinting =
_config.getProperty(Broker.CONFIG_OF_BROKER_ENABLE_QUERY_FINGERPRINTING,
+ boolean fingerprintingConfigured =
_config.getProperty(Broker.CONFIG_OF_BROKER_ENABLE_QUERY_FINGERPRINTING,
Broker.DEFAULT_BROKER_ENABLE_QUERY_FINGERPRINTING);
+ boolean redactionNeedsFingerprinting =
+ _queryLogger.getSqlRedactionMode() ==
QueryLogger.SqlRedactionMode.LITERAL_VALUES;
+ if (redactionNeedsFingerprinting && !fingerprintingConfigured) {
+ LOGGER.warn("SQL redaction mode 'literal_values' requires query
fingerprinting. "
+ + "Enabling query fingerprinting automatically.");
+ }
+ _enableQueryFingerprinting = fingerprintingConfigured ||
redactionNeedsFingerprinting;
if (_enableMultistageMigrationMetric) {
_multistageCompileExecutor = Executors.newSingleThreadExecutor();
_multistageCompileQueryQueue = new LinkedBlockingQueue<>(1000);
@@ -299,7 +306,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
// we can get the cid from QueryThreadContext
serverUrls.add(Pair.of(String.format("%s/query/%s",
serverInstance.getAdminEndpoint(), globalQueryId), null));
}
- LOGGER.debug("Cancelling the query: {} via server urls: {}",
queryServers._query, serverUrls);
+ LOGGER.debug("Cancelling the query: {} via server urls: {}",
_queryLogger.redactQuery(queryServers._query),
+ serverUrls);
CompletionService<MultiHttpRequestResponse> completionService =
new MultiHttpRequest(executor, connMgr).execute(serverUrls, null,
timeoutMs, "DELETE", HttpDelete::new);
List<String> errMsgs = new ArrayList<>(serverUrls.size());
@@ -321,7 +329,7 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
serverResponses.put(uri.getHost() + ":" + uri.getPort(), status);
}
} catch (Exception e) {
- LOGGER.error("Failed to cancel query: {}", queryServers._query, e);
+ LOGGER.error("Failed to cancel query: {}",
_queryLogger.redactQuery(queryServers._query), e);
// Can't just throw exception from here as there is a need to release
the other connections.
// So just collect the error msg to throw them together after the
for-loop.
errMsgs.add(e.getMessage());
@@ -342,21 +350,23 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
JsonNode request, @Nullable RequesterIdentity requesterIdentity,
RequestContext requestContext,
@Nullable HttpHeaders httpHeaders, AccessControl accessControl)
throws Exception {
- boolean queryWasLogged = _queryLogger.logQueryReceived(requestId, query);
-
+ QueryFingerprint queryFingerprint = null;
String queryHash = CommonConstants.Broker.DEFAULT_QUERY_HASH;
if (_enableQueryFingerprinting) {
try {
- QueryFingerprint queryFingerprint =
QueryFingerprintUtils.generateFingerprint(sqlNodeAndOptions);
+ queryFingerprint =
QueryFingerprintUtils.generateFingerprint(sqlNodeAndOptions);
if (queryFingerprint != null) {
queryHash = queryFingerprint.getQueryHash();
requestContext.setQueryFingerprint(queryFingerprint);
}
} catch (Exception e) {
- LOGGER.warn("Failed to generate query fingerprint for request {}: {}.
{}", requestId, query, e.getMessage());
+ LOGGER.warn("Failed to generate query fingerprint for request {}: {}.
{}", requestId,
+ _queryLogger.redactQuery(query), e.getMessage());
}
}
+ boolean queryWasLogged = _queryLogger.logQueryReceived(requestId, query,
queryFingerprint);
+
String cid = extractClientRequestId(sqlNodeAndOptions);
if (cid == null) {
cid = Long.toString(requestId);
@@ -487,14 +497,16 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
// Validate QPS
if (!_queryQuotaManager.acquireDatabase(database)) {
String errorMessage =
- String.format("Request %d: %s exceeds query quota for database:
%s", requestId, query, database);
+ String.format("Request %d: %s exceeds query quota for database:
%s", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), database);
LOGGER.info(errorMessage);
requestContext.setErrorCode(QueryErrorCode.TOO_MANY_REQUESTS);
return new BrokerResponseNative(QueryErrorCode.TOO_MANY_REQUESTS,
errorMessage);
}
if (!_queryQuotaManager.acquireLogicalTable(tableName)) {
String errorMessage =
- String.format("Request %d: %s exceeds query quota for table: %s.",
requestId, query, tableName);
+ String.format("Request %d: %s exceeds query quota for table: %s.",
requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), tableName);
requestContext.setErrorCode(QueryErrorCode.TOO_MANY_REQUESTS);
return new BrokerResponseNative(QueryErrorCode.TOO_MANY_REQUESTS,
errorMessage);
}
@@ -544,14 +556,16 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
// Validate QPS quota
if (!_queryQuotaManager.acquireDatabase(database)) {
String errorMessage =
- String.format("Request %d: %s exceeds query quota for database:
%s", requestId, query, database);
+ String.format("Request %d: %s exceeds query quota for database:
%s", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), database);
LOGGER.info(errorMessage);
requestContext.setErrorCode(QueryErrorCode.TOO_MANY_REQUESTS);
return new BrokerResponseNative(QueryErrorCode.TOO_MANY_REQUESTS,
errorMessage);
}
if (!_queryQuotaManager.acquire(tableName)) {
String errorMessage =
- String.format("Request %d: %s exceeds query quota for table: %s",
requestId, query, tableName);
+ String.format("Request %d: %s exceeds query quota for table: %s",
requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), tableName);
LOGGER.info(errorMessage);
requestContext.setErrorCode(QueryErrorCode.TOO_MANY_REQUESTS);
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.QUERY_QUOTA_EXCEEDED, 1);
@@ -602,13 +616,15 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
TableRouteInfo routeInfo = routeProvider.getTableRouteInfo(tableName,
_tableCache, selectedRoutingManager);
if (!routeInfo.isExists()) {
- LOGGER.info("Table not found for request {}: {}", requestId, query);
+ LOGGER.info("Table not found for request {}: {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.TABLE_DOES_NOT_EXIST);
return BrokerResponseNative.TABLE_DOES_NOT_EXIST;
}
if (!routeInfo.isRouteExists()) {
- LOGGER.info("No table matches for request {}: {}", requestId, query);
+ LOGGER.info("No table matches for request {}: {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.BROKER_RESOURCE_MISSING);
_brokerMetrics.addMeteredGlobalValue(BrokerMeter.RESOURCE_MISSING_EXCEPTIONS,
1);
return BrokerResponseNative.NO_TABLE_RESULT;
@@ -632,7 +648,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
try {
validateRequest(serverPinotQuery, _queryResponseLimit);
} catch (Exception e) {
- LOGGER.info("Caught exception while validating request {}: {}, {}",
requestId, query, e.getMessage());
+ LOGGER.info("Caught exception while validating request {}: {}, {}",
requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
requestContext.setErrorCode(QueryErrorCode.QUERY_VALIDATION);
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.QUERY_VALIDATION_EXCEPTIONS, 1);
return new BrokerResponseNative(QueryErrorCode.QUERY_VALIDATION,
e.getMessage());
@@ -648,7 +665,7 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
// Attempt to add the query to the compile queue; drop if queue is full
if (!_multistageCompileQueryQueue.offer(Pair.of(query, database))) {
LOGGER.trace("Not compiling query `{}` using the multi-stage query
engine because the query queue is full",
- query);
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
}
}
@@ -755,7 +772,7 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
if (disabledTableNames != null) {
for (String name : disabledTableNames) {
String errorMessage = String.format("%s Table is disabled", name);
- LOGGER.info("{}: {}", errorMessage, query);
+ LOGGER.info("{}: {}", errorMessage, _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
errorMsgs.add(new
QueryProcessingException(QueryErrorCode.TABLE_IS_DISABLED, errorMessage));
}
}
@@ -784,7 +801,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
QueryProcessingException firstErrorMsg = errorMsgs.get(0);
String logTail = errorMsgs.size() > 1 ? (errorMsgs.size()) + "
errorMsgs found. Logging only the first one"
: "1 exception found";
- LOGGER.info("No server found for request {}: {}. {} {}", requestId,
query, logTail, firstErrorMsg);
+ LOGGER.info("No server found for request {}: {}. {} {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), logTail, firstErrorMsg);
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.NO_SERVER_FOUND_EXCEPTIONS, 1);
return BrokerResponseNative.fromBrokerErrors(errorMsgs);
} else {
@@ -823,7 +841,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
}
} catch (TimeoutException e) {
String errorMessage = e.getMessage();
- LOGGER.info("{} {}: {}", errorMessage, requestId, query);
+ LOGGER.info("{} {}: {}", errorMessage, requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.REQUEST_TIMEOUT_BEFORE_SCATTERED_EXCEPTIONS, 1);
errorMsgs.add(new
QueryProcessingException(QueryErrorCode.BROKER_TIMEOUT, errorMessage));
return BrokerResponseNative.fromBrokerErrors(errorMsgs);
@@ -1048,7 +1067,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
try {
pinotQuery = CalciteSqlParser.compileToPinotQuery(sqlNodeAndOptions);
} catch (Exception e) {
- LOGGER.info("Caught exception while compiling SQL request {}: {}, {}",
requestId, query, e.getMessage());
+ LOGGER.info("Caught exception while compiling SQL request {}: {}, {}",
requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
_brokerMetrics.addMeteredGlobalValue(BrokerMeter.REQUEST_COMPILATION_EXCEPTIONS,
1);
requestContext.setErrorCode(QueryErrorCode.SQL_PARSING);
// Check if the query is a v2 supported query
@@ -1080,7 +1100,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
}
if (isLiteralOnlyQuery(pinotQuery)) {
- LOGGER.debug("Request {} contains only Literal, skipping server query:
{}", requestId, query);
+ LOGGER.debug("Request {} contains only Literal, skipping server query:
{}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
try {
if (pinotQuery.isExplain()) {
// EXPLAIN PLAN results to show that query is evaluated exclusively
by Broker.
@@ -1090,25 +1111,28 @@ public abstract class
BaseSingleStageBrokerRequestHandler extends BaseBrokerRequ
} catch (Exception e) {
// TODO: refine the exceptions here to early termination the queries
won't requires to send to servers.
LOGGER.warn("Unable to execute literal request {}: {} at broker,
fallback to server query. {}", requestId,
- query, e.getMessage());
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
}
}
PinotQuery serverPinotQuery = GapfillUtils.stripGapfill(pinotQuery);
DataSource dataSource = serverPinotQuery.getDataSource();
if (dataSource == null) {
- LOGGER.info("Data source (FROM clause) not found in request {}: {}",
requestId, query);
+ LOGGER.info("Data source (FROM clause) not found in request {}: {}",
requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.QUERY_VALIDATION);
return new CompileResult(new BrokerResponseNative(
QueryErrorCode.QUERY_VALIDATION, "Data source (FROM clause) not
found"));
}
if (dataSource.getJoin() != null) {
- LOGGER.info("JOIN is not supported in request {}: {}", requestId, query);
+ LOGGER.info("JOIN is not supported in request {}: {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.QUERY_VALIDATION);
return new CompileResult(new
BrokerResponseNative(QueryErrorCode.QUERY_VALIDATION, "JOIN is not supported"));
}
if (dataSource.getTableName() == null) {
- LOGGER.info("Table name not found in request {}: {}", requestId, query);
+ LOGGER.info("Table name not found in request {}: {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.QUERY_VALIDATION);
return new CompileResult(new
BrokerResponseNative(QueryErrorCode.QUERY_VALIDATION, "Table name not found"));
}
@@ -1117,8 +1141,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
handleSubquery(serverPinotQuery, requestId, request, requesterIdentity,
requestContext, httpHeaders,
accessControl);
} catch (Exception e) {
- LOGGER.info("Caught exception while handling the subquery in request {}:
{}, {}", requestId, query,
- e.getMessage());
+ LOGGER.info("Caught exception while handling the subquery in request {}:
{}, {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
requestContext.setErrorCode(QueryErrorCode.QUERY_EXECUTION);
return new CompileResult(
new BrokerResponseNative(QueryErrorCode.QUERY_EXECUTION,
e.getMessage()));
@@ -1131,7 +1155,8 @@ public abstract class BaseSingleStageBrokerRequestHandler
extends BaseBrokerRequ
getActualTableName(DatabaseUtils.translateTableName(dataSource.getTableName(),
httpHeaders, ignoreCase),
_tableCache);
} catch (DatabaseConflictException e) {
- LOGGER.info("{}. Request {}: {}", e.getMessage(), requestId, query);
+ LOGGER.info("{}. Request {}: {}", e.getMessage(), requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()));
_brokerMetrics.addMeteredGlobalValue(BrokerMeter.QUERY_VALIDATION_EXCEPTIONS,
1);
requestContext.setErrorCode(QueryErrorCode.QUERY_VALIDATION);
return new CompileResult(
@@ -1156,15 +1181,15 @@ public abstract class
BaseSingleStageBrokerRequestHandler extends BaseBrokerRequ
} catch (Exception e) {
// Throw exceptions with column in-existence error.
if (e instanceof BadQueryRequestException) {
- LOGGER.info("Caught exception while checking column names in request
{}: {}, {}", requestId, query,
- e.getMessage());
+ LOGGER.info("Caught exception while checking column names in request
{}: {}, {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
requestContext.setErrorCode(QueryErrorCode.UNKNOWN_COLUMN);
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.UNKNOWN_COLUMN_EXCEPTIONS, 1);
return new CompileResult(
new BrokerResponseNative(QueryErrorCode.UNKNOWN_COLUMN,
e.getMessage()));
}
- LOGGER.warn("Caught exception while updating column names in request {}:
{}, {}", requestId, query,
- e.getMessage());
+ LOGGER.warn("Caught exception while updating column names in request {}:
{}, {}", requestId,
+ _queryLogger.redactQuery(query,
requestContext.getQueryFingerprint()), e.getMessage());
}
if (_defaultHllLog2m > 0) {
@@ -1357,13 +1382,15 @@ public abstract class
BaseSingleStageBrokerRequestHandler extends BaseBrokerRequ
serverPinotQuery.setQueryOptions(queryOptions);
CalciteSqlParser.queryRewrite(serverPinotQuery, RlsFiltersRewriter.class);
_brokerMetrics.addMeteredTableValue(rawTableName,
BrokerMeter.RLS_FILTERS_APPLIED, 1);
- LOGGER.debug("Applied RLS filters for request {} on table {}: {}",
requestId, rawTableName, query);
+ LOGGER.debug("Applied RLS filters for request {} on table {}: {}",
requestId, rawTableName,
+ _queryLogger.redactQuery(query));
}
private void throwAccessDeniedError(long requestId, String query,
RequestContext requestContext, String tableName,
AuthorizationResult authorizationResult) {
_brokerMetrics.addMeteredTableValue(tableName,
BrokerMeter.REQUEST_DROPPED_DUE_TO_ACCESS_ERROR, 1);
- LOGGER.info("Access denied for request {}: {}, table: {}, reason :{}",
requestId, query, tableName,
+ LOGGER.info("Access denied for request {}: {}, table: {}, reason :{}",
requestId,
+ _queryLogger.redactQuery(query, requestContext.getQueryFingerprint()),
tableName,
authorizationResult.getFailureMessage());
requestContext.setErrorCode(QueryErrorCode.ACCESS_DENIED);
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
index d9d60bbaad4..91e2d2c2c9a 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
@@ -231,9 +231,16 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
_config.containsKey(CommonConstants.Broker.CONFIG_OF_BROKER_MSE_PLANNER_DISABLED_RULES)
? Set.copyOf(
_config.getProperty(CommonConstants.Broker.CONFIG_OF_BROKER_MSE_PLANNER_DISABLED_RULES,
List.of()))
: CommonConstants.Broker.DEFAULT_DISABLED_RULES;
- _enableQueryFingerprinting = _config.getProperty(
+ boolean fingerprintingConfigured = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_BROKER_ENABLE_QUERY_FINGERPRINTING,
CommonConstants.Broker.DEFAULT_BROKER_ENABLE_QUERY_FINGERPRINTING);
+ boolean redactionNeedsFingerprinting =
+ _queryLogger.getSqlRedactionMode() ==
QueryLogger.SqlRedactionMode.LITERAL_VALUES;
+ if (redactionNeedsFingerprinting && !fingerprintingConfigured) {
+ LOGGER.warn("SQL redaction mode 'literal_values' requires query
fingerprinting. "
+ + "Enabling query fingerprinting automatically.");
+ }
+ _enableQueryFingerprinting = fingerprintingConfigured ||
redactionNeedsFingerprinting;
int streamingGroupByFlushThreshold = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_MSE_STREAMING_GROUP_BY_FLUSH_THRESHOLD,
CommonConstants.Broker.DEFAULT_MSE_STREAMING_GROUP_BY_FLUSH_THRESHOLD);
@@ -394,21 +401,23 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
protected BrokerResponse handleRequestThrowing(long requestId, String query,
SqlNodeAndOptions sqlNodeAndOptions,
@Nullable RequesterIdentity requesterIdentity, RequestContext
requestContext, HttpHeaders httpHeaders)
throws QueryException, WebApplicationException {
- boolean queryWasLogged = _queryLogger.logQueryReceived(requestId, query);
-
+ QueryFingerprint queryFingerprint = null;
String queryHash = CommonConstants.Broker.DEFAULT_QUERY_HASH;
if (_enableQueryFingerprinting) {
try {
- QueryFingerprint queryFingerprint =
QueryFingerprintUtils.generateFingerprint(sqlNodeAndOptions);
+ queryFingerprint =
QueryFingerprintUtils.generateFingerprint(sqlNodeAndOptions);
if (queryFingerprint != null) {
queryHash = queryFingerprint.getQueryHash();
requestContext.setQueryFingerprint(queryFingerprint);
}
} catch (Exception e) {
- LOGGER.warn("Failed to generate query fingerprint for request {}: {}.
{}", requestId, query, e.getMessage());
+ LOGGER.warn("Failed to generate query fingerprint for request {}: {}.
{}", requestId,
+ _queryLogger.redactQuery(query), e.getMessage());
}
}
+ boolean queryWasLogged = _queryLogger.logQueryReceived(requestId, query,
queryFingerprint);
+
String cid = extractClientRequestId(sqlNodeAndOptions);
if (cid == null) {
cid = Long.toString(requestId);
@@ -702,7 +711,8 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
// Validate QPS quota
if (hasExceededQPSQuota(query.getDatabase(), tableNames, requestContext)) {
- String errorMessage = String.format("Request %d: %s exceeds query
quota.", requestId, query);
+ String errorMessage = String.format("Request %d: %s exceeds query
quota.", requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()));
return new BrokerResponseNative(QueryErrorCode.TOO_MANY_REQUESTS,
errorMessage);
}
@@ -715,14 +725,16 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
// these requests.
if (!_queryThrottler.tryAcquire(estimatedNumQueryThreads,
timer.getRemainingTimeMs(),
TimeUnit.MILLISECONDS)) {
- LOGGER.warn("Timed out waiting to execute request {}: {}", requestId,
query);
+ LOGGER.warn("Timed out waiting to execute request {}: {}", requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.EXECUTION_TIMEOUT);
return new BrokerResponseNative(QueryErrorCode.EXECUTION_TIMEOUT);
}
_brokerMetrics.setValueOfGlobalGauge(BrokerGauge.ESTIMATED_MSE_SERVER_THREADS,
_queryThrottler.currentQueryServerThreads());
} catch (InterruptedException e) {
- LOGGER.warn("Interrupt received while waiting to execute request {}:
{}", requestId, query);
+ LOGGER.warn("Interrupt received while waiting to execute request {}:
{}", requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()));
requestContext.setErrorCode(QueryErrorCode.EXECUTION_TIMEOUT);
return new BrokerResponseNative(QueryErrorCode.EXECUTION_TIMEOUT);
}
@@ -746,7 +758,8 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
} catch (Throwable t) {
QueryErrorCode queryErrorCode = QueryErrorCode.QUERY_EXECUTION;
String consolidatedMessage =
ExceptionUtils.consolidateExceptionTraces(t);
- LOGGER.error("Caught exception reducing all-leaf-empty request {}:
{}, {}", requestId, query,
+ LOGGER.error("Caught exception reducing all-leaf-empty request {}:
{}, {}", requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()),
consolidatedMessage);
requestContext.setErrorCode(queryErrorCode);
return new BrokerResponseNative(queryErrorCode, consolidatedMessage);
@@ -768,7 +781,9 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
} catch (Throwable t) {
QueryErrorCode queryErrorCode = QueryErrorCode.QUERY_EXECUTION;
String consolidatedMessage =
ExceptionUtils.consolidateExceptionTraces(t);
- LOGGER.error("Caught exception executing request {}: {}, {}",
requestId, query, consolidatedMessage);
+ LOGGER.error("Caught exception executing request {}: {}, {}",
requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()),
+ consolidatedMessage);
requestContext.setErrorCode(queryErrorCode);
return new BrokerResponseNative(queryErrorCode, consolidatedMessage);
} finally {
@@ -788,7 +803,8 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
for (String table : tableNames) {
_brokerMetrics.addMeteredTableValue(table,
BrokerMeter.BROKER_RESPONSES_WITH_TIMEOUTS, 1);
}
- LOGGER.warn("Timed out executing request {}: {}", requestId, query);
+ LOGGER.warn("Timed out executing request {}: {}", requestId,
+ _queryLogger.redactQuery(query.getTextQuery(),
requestContext.getQueryFingerprint()));
}
requestContext.setErrorCode(errorCode);
} else {
@@ -921,17 +937,17 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
return queryPlanResultFuture.get(timer.getRemainingTimeMs(),
TimeUnit.MILLISECONDS);
} catch (TimeoutException e) {
String errorMsg = "Timed out while planning query";
- LOGGER.warn(errorMsg + " {}", query, e);
+ LOGGER.warn(errorMsg + " {}", _queryLogger.redactQuery(query), e);
queryPlanResultFuture.cancel(true);
throw QueryErrorCode.BROKER_TIMEOUT.asException(errorMsg);
} catch (InterruptedException e) {
- LOGGER.warn("Interrupt received while planning query {}: {}", requestId,
query);
+ LOGGER.warn("Interrupt received while planning query {}: {}", requestId,
_queryLogger.redactQuery(query));
throw QueryErrorCode.INTERNAL.asException("Interrupted while planning
query");
} catch (ExecutionException e) {
if (e.getCause() instanceof QueryException) {
throw (QueryException) e.getCause();
} else {
- LOGGER.warn("Error while planning query {}: {}", query, e.getCause());
+ LOGGER.warn("Error while planning query {}: {}",
_queryLogger.redactQuery(query), e.getCause());
throw QueryErrorCode.INTERNAL.asException("Error while planning
query", e.getCause());
}
}
diff --git
a/pinot-broker/src/test/java/org/apache/pinot/broker/querylog/QueryLoggerTest.java
b/pinot-broker/src/test/java/org/apache/pinot/broker/querylog/QueryLoggerTest.java
index fe53d5e0016..80583db3b5f 100644
---
a/pinot-broker/src/test/java/org/apache/pinot/broker/querylog/QueryLoggerTest.java
+++
b/pinot-broker/src/test/java/org/apache/pinot/broker/querylog/QueryLoggerTest.java
@@ -43,6 +43,7 @@ import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
+import static org.apache.pinot.broker.querylog.QueryLogger.SqlRedactionMode;
import static org.mockito.MockitoAnnotations.openMocks;
@@ -98,7 +99,8 @@ public class QueryLoggerTest {
public void shouldFormatLogLineProperly() {
// Given:
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, true);
@@ -139,7 +141,8 @@ public class QueryLoggerTest {
public void shouldOmitClientId() {
// Given:
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, false,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, false,
true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, true);
@@ -154,7 +157,8 @@ public class QueryLoggerTest {
public void shouldNotLogCompletionWhenWasLoggedFalseAndNoForceLog() {
// Given: wasLogged=false and no force-log conditions (no exceptions, not
slow)
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, false);
@@ -167,7 +171,8 @@ public class QueryLoggerTest {
public void shouldForceLogWhenNumGroupsLimitIsReached() {
// Given: wasLogged=false but numGroupsLimitReached (force-log condition)
QueryLogger.QueryLogParams params = generateParams(true, true, 0, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, false);
@@ -180,7 +185,8 @@ public class QueryLoggerTest {
public void shouldForceLogWhenExceptionsExist() {
// Given: wasLogged=false but exceptions exist (force-log condition)
QueryLogger.QueryLogParams params = generateParams(false, false, 1, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, false);
@@ -193,7 +199,8 @@ public class QueryLoggerTest {
public void shouldForceLogWhenTimeIsMoreThanOneSecond() {
// Given: wasLogged=false but query took >1s (force-log condition)
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 1456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, false);
@@ -206,10 +213,11 @@ public class QueryLoggerTest {
public void shouldLogQueryReceivedWhenAllowed() {
// Given: rate limiter allows
Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(true);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
- boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo");
+ boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo", null);
// Then:
Assert.assertTrue(wasLogged);
@@ -221,10 +229,11 @@ public class QueryLoggerTest {
public void shouldNotLogQueryReceivedWhenRateLimited() {
// Given: rate limiter denies
Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(false);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
- boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo");
+ boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo", null);
// Then:
Assert.assertFalse(wasLogged);
@@ -235,10 +244,11 @@ public class QueryLoggerTest {
public void shouldReturnTrueButNotLogWhenLogBeforeProcessingIsDisabled() {
// Given: rate limiter allows, but logBeforeProcessing=false
Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(true);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
false, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
false,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
- boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo");
+ boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM
foo", null);
// Then: returns true because rate limiter allowed, but no log because
logBeforeProcessing=false
Assert.assertTrue(wasLogged);
@@ -249,7 +259,8 @@ public class QueryLoggerTest {
public void shouldLogCompletionWhenWasLoggedIsTrue() {
// Given:
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, true);
@@ -279,22 +290,25 @@ public class QueryLoggerTest {
return true;
});
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
ExecutorService executorService = Executors.newSingleThreadExecutor();
// When:
try {
- Assert.assertFalse(queryLogger.logQueryReceived(123, "SELECT * FROM
foo")); // 1 this one gets dropped
+ Assert.assertFalse(queryLogger.logQueryReceived(123, "SELECT * FROM
foo", null)); // 1 this one gets dropped
// 2 this one succeeds, but blocks when it checks whether to log the
dropped count
Future<Boolean> blockedLogger =
- executorService.submit(() -> queryLogger.logQueryReceived(123,
"SELECT * FROM foo"));
+ executorService.submit(() -> queryLogger.logQueryReceived(123,
"SELECT * FROM foo", null));
Assert.assertTrue(firstDroppedLogAttempted.await(5, TimeUnit.SECONDS),
"expected the first successful log to reach the dropped-log rate
limiter");
- Assert.assertFalse(queryLogger.logQueryReceived(123, "SELECT * FROM
foo")); // 3 this one gets dropped
- Assert.assertTrue(queryLogger.logQueryReceived(123, "SELECT * FROM
foo")); // 4 this one drains the dropped count
+ // 3 this one gets dropped
+ Assert.assertFalse(queryLogger.logQueryReceived(123, "SELECT * FROM
foo", null));
+ // 4 this one drains the dropped count
+ Assert.assertTrue(queryLogger.logQueryReceived(123, "SELECT * FROM foo",
null));
releaseFirstDroppedLogAttempt.countDown();
Assert.assertTrue(blockedLogger.get(5, TimeUnit.SECONDS));
@@ -314,7 +328,8 @@ public class QueryLoggerTest {
// Given:
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
new QueryFingerprint("abc", "SELECT * FROM foo"));
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, true);
@@ -330,7 +345,8 @@ public class QueryLoggerTest {
public void shouldEmitEmptyQueryHashWhenNotSet() {
// Given:
QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
- QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true,
true, _logger, _droppedRateLimiter);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.NONE, _logger, _droppedRateLimiter);
// When:
queryLogger.logQueryCompleted(params, true);
@@ -342,6 +358,117 @@ public class QueryLoggerTest {
"Expected empty queryHash field. Got: " + logLine);
}
+ @Test
+ public void shouldRedactQueryInLogQueryReceivedWhenEnabled() {
+ // Given:
+ Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(true);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.LITERAL_VALUES, _logger, _droppedRateLimiter);
+ QueryFingerprint fingerprint = new QueryFingerprint("abc", "SELECT * FROM
foo WHERE id = ?");
+
+ // When:
+ boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM foo
WHERE id = 42", fingerprint);
+
+ // Then:
+ Assert.assertTrue(wasLogged);
+ Assert.assertEquals(_infoLog.size(), 1);
+ Assert.assertTrue(_infoLog.get(0).contains("SELECT * FROM foo WHERE id =
?"),
+ "Expected redacted query. Got: " + _infoLog.get(0));
+ Assert.assertFalse(_infoLog.get(0).contains("42"),
+ "Raw literal should not appear. Got: " + _infoLog.get(0));
+ }
+
+ @Test
+ public void
shouldLogSentinelInReceivedWhenFingerprintNullAndRedactionEnabled() {
+ // Given:
+ Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(true);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.LITERAL_VALUES, _logger, _droppedRateLimiter);
+
+ // When:
+ boolean wasLogged = queryLogger.logQueryReceived(123L, "SELECT * FROM foo
WHERE id = 42", null);
+
+ // Then:
+ Assert.assertTrue(wasLogged);
+ Assert.assertEquals(_infoLog.size(), 1);
+
Assert.assertTrue(_infoLog.get(0).contains("FINGERPRINT_FAILED_QUERY_REDACTED"),
+ "Expected sentinel. Got: " + _infoLog.get(0));
+ Assert.assertFalse(_infoLog.get(0).contains("42"),
+ "Raw literal should not appear. Got: " + _infoLog.get(0));
+ }
+
+ @Test
+ public void shouldRedactQueryInLogQueryCompletedWhenEnabled() {
+ // Given:
+ QueryFingerprint fingerprint = new QueryFingerprint("abc", "SELECT * FROM
foo WHERE id = ?");
+ QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
fingerprint);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.LITERAL_VALUES, _logger, _droppedRateLimiter);
+
+ // When:
+ queryLogger.logQueryCompleted(params, true);
+
+ // Then:
+ Assert.assertEquals(_infoLog.size(), 1);
+ String logLine = _infoLog.get(0);
+ Assert.assertTrue(logLine.contains("query=SELECT * FROM foo WHERE id = ?"),
+ "Expected redacted query in completion log. Got: " + logLine);
+ }
+
+ @Test
+ public void
shouldLogSentinelInCompletedWhenFingerprintNullAndRedactionEnabled() {
+ // Given:
+ QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.LITERAL_VALUES, _logger, _droppedRateLimiter);
+
+ // When:
+ queryLogger.logQueryCompleted(params, true);
+
+ // Then:
+ Assert.assertEquals(_infoLog.size(), 1);
+ String logLine = _infoLog.get(0);
+
Assert.assertTrue(logLine.contains("query=FINGERPRINT_FAILED_QUERY_REDACTED"),
+ "Expected sentinel in completion log. Got: " + logLine);
+ }
+
+ @Test
+ public void shouldFullyRedactQueryInLogQueryReceived() {
+ // Given:
+ Mockito.when(_logRateLimiter.tryAcquire()).thenReturn(true);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.FULL, _logger, _droppedRateLimiter);
+
+ // When:
+ queryLogger.logQueryReceived(123L, "SELECT * FROM foo WHERE id = 42",
null);
+
+ // Then:
+ Assert.assertEquals(_infoLog.size(), 1);
+ Assert.assertTrue(_infoLog.get(0).contains("REDACTED"),
+ "Expected REDACTED. Got: " + _infoLog.get(0));
+ Assert.assertFalse(_infoLog.get(0).contains("foo"),
+ "No part of SQL should appear. Got: " + _infoLog.get(0));
+ }
+
+ @Test
+ public void shouldFullyRedactQueryInLogQueryCompleted() {
+ // Given:
+ QueryLogger.QueryLogParams params = generateParams(false, false, 0, 456,
null);
+ QueryLogger queryLogger = new QueryLogger(_logRateLimiter, 100, true, true,
+ SqlRedactionMode.FULL, _logger, _droppedRateLimiter);
+
+ // When:
+ queryLogger.logQueryCompleted(params, true);
+
+ // Then:
+ Assert.assertEquals(_infoLog.size(), 1);
+ String logLine = _infoLog.get(0);
+ Assert.assertTrue(logLine.contains("query=REDACTED"),
+ "Expected fully redacted query. Got: " + logLine);
+ Assert.assertFalse(logLine.contains("SELECT"),
+ "No SQL should appear. Got: " + logLine);
+ }
+
private QueryLogger.QueryLogParams generateParams(boolean
numGroupsLimitReached, boolean numGroupsWarningLimitReached,
int numExceptions, long timeUsedMs) {
return generateParams(numGroupsLimitReached, numGroupsWarningLimitReached,
numExceptions, timeUsedMs, null);
diff --git
a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
index 1c8dca6fdb1..3ddfaf79c0f 100644
--- a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
+++ b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
@@ -381,6 +381,9 @@ public class CommonConstants {
public static final String CONFIG_OF_BROKER_QUERY_LOG_BEFORE_PROCESSING =
"pinot.broker.query.log.logBeforeProcessing";
public static final boolean DEFAULT_BROKER_QUERY_LOG_BEFORE_PROCESSING =
true;
+ public static final String CONFIG_OF_BROKER_QUERY_LOG_SQL_REDACTION =
+ "pinot.broker.query.log.sqlRedaction";
+ public static final String DEFAULT_BROKER_QUERY_LOG_SQL_REDACTION = "none";
public static final String CONFIG_OF_BROKER_QUERY_ENABLE_NULL_HANDLING =
"pinot.broker.query.enable.null.handling";
/**
* When true, the broker initializes the materialized view metadata cache
and query rewrite
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]