FrankChen021 commented on code in PR #18525:
URL: https://github.com/apache/druid/pull/18525#discussion_r3957998667
##########
indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/common/OrderedPartitionableRecord.java:
##########
@@ -50,15 +51,16 @@ public OrderedPartitionableRecord(
List<RecordType> data
)
{
- this(stream, partitionId, sequenceNumber, data, null);
+ this(stream, partitionId, sequenceNumber, data, null, false);
}
public OrderedPartitionableRecord(
String stream,
PartitionIdType partitionId,
SequenceOffsetType sequenceNumber,
List<RecordType> data,
- Long timestamp
+ Long timestamp,
+ boolean filtered
Review Comment:
[P1] Preserve the public five-argument constructor
This changes the public `(String, PartitionIdType, SequenceOffsetType,
List<RecordType>, Long)` constructor into a six-argument signature. Existing
stream extensions compiled against this class will fail with
`NoSuchMethodError`, and source extensions using the constructor will no longer
compile. Keep the old overload and delegate it to the new one with `filtered =
false`.
##########
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaSamplerSpec.java:
##########
@@ -74,7 +74,9 @@ protected KafkaRecordSupplier createRecordSupplier()
objectMapper,
kafkaSupervisorIOConfig.getConfigOverrides(),
kafkaSupervisorIOConfig.isMultiTopic(),
- null
+ null,
+ // Apply the same header-based filter during sampling so sampled
records match what ingestion tasks keep.
+ kafkaSupervisorIOConfig.getheaderBasedFilterConfig()
Review Comment:
[P2] Skip filtered records in Kafka sampling
The filter is still ineffective for sampling. `SeekableStreamSamplerSpec`
wraps this supplier in `RecordSupplierInputSource`, whose iterator consumes
`recordIterator.next().getData()` without checking
`OrderedPartitionableRecord.isFiltered()`. Because `KafkaRecordSupplier`
intentionally retains the filtered payload in the marker, sample queries will
parse and return records that ingestion tasks discard. Make the sampler input
source skip filtered records while retaining payloads for task byte accounting,
or use a sampler-specific record path.
##########
extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaHeaderBasedFilterEvaluator.java:
##########
@@ -0,0 +1,182 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.indexing.kafka;
+
+import com.github.benmanes.caffeine.cache.Cache;
+import com.github.benmanes.caffeine.cache.Caffeine;
+import org.apache.druid.indexing.kafka.supervisor.KafkaHeaderBasedFilterConfig;
+import org.apache.druid.java.util.common.logger.Logger;
+import org.apache.druid.query.filter.Filter;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.common.header.Header;
+import org.apache.kafka.common.header.Headers;
+
+import javax.annotation.Nullable;
+
+import java.io.UncheckedIOException;
+import java.nio.ByteBuffer;
+import java.nio.charset.CharacterCodingException;
+import java.nio.charset.Charset;
+import java.nio.charset.CharsetDecoder;
+import java.nio.charset.CodingErrorAction;
+
+/**
+ * Evaluates Kafka header filters for pre-ingestion filtering.
+ */
+public class KafkaHeaderBasedFilterEvaluator
+{
+ private static final Logger log = new
Logger(KafkaHeaderBasedFilterEvaluator.class);
+
+ /**
+ * Header values larger than this are decoded but not cached, so that a
bounded number of cache entries
+ * (see {@link KafkaHeaderBasedFilterConfig#getStringDecodingCacheSize()})
cannot retain unbounded heap when
+ * distinct large headers are seen. Total cached bytes are therefore bounded
by cacheSize * this cap.
+ */
+ private static final int MAX_CACHEABLE_HEADER_BYTES = 4096;
+
+ private final HeaderFilterHandler filterHandler;
+ private final String headerName;
+ private final Charset encoding;
+ private final Cache<ByteBuffer, String> stringDecodingCache;
+
+ /**
+ * Creates a new KafkaHeaderBasedFilterEvaluator with the given
configuration.
+ *
+ * @param headerBasedFilterConfig the configuration containing filter,
encoding, and cache settings
+ * @throws IllegalArgumentException if the filter type is not supported
+ */
+ public KafkaHeaderBasedFilterEvaluator(KafkaHeaderBasedFilterConfig
headerBasedFilterConfig)
+ {
+ this.encoding = Charset.forName(headerBasedFilterConfig.getEncoding());
+ this.stringDecodingCache = Caffeine.newBuilder()
+ .maximumSize(headerBasedFilterConfig.getStringDecodingCacheSize())
+ .build();
+
+ Filter filter = headerBasedFilterConfig.getFilter().toFilter();
+ this.filterHandler = HeaderFilterHandlerFactory.forFilter(filter);
+ this.headerName = filterHandler.getHeaderName();
+
+ log.info("Initialized Kafka header filter: %s with encoding [%s] and cache
size [%d]",
+ filterHandler.getDescription(),
+ headerBasedFilterConfig.getEncoding(),
+ headerBasedFilterConfig.getStringDecodingCacheSize());
+ }
+
+
+ /**
+ * Evaluates whether a Kafka record should be included based on its headers.
+ *
+ * @param record the Kafka consumer record
+ * @return true if the record should be included, false if it should be
filtered out
+ */
+ public boolean shouldIncludeRecord(ConsumerRecord<byte[], byte[]> record)
+ {
+ try {
+ return evaluateInclusion(record.headers());
+ }
+ catch (Exception e) {
+ log.warn(
+ e,
+ "Error evaluating header filter for record at topic [%s] partition
[%d] offset [%d], including record",
+ record.topic(),
+ record.partition(),
+ record.offset()
+ );
+ return true; // Default to including record on error
+ }
+ }
+
+ /**
+ * Evaluates whether a record should be included based on its headers.
+ *
+ * Uses permissive behavior: records with missing, null, or undecodable
headers
+ * are included by default. Only records with successfully decoded header
values
+ * that don't match the filter criteria are excluded.
+ *
+ * @param headers the Kafka message headers to evaluate
+ * @return true if the record should be included, false if it should be
filtered out
+ */
+ private boolean evaluateInclusion(Headers headers)
+ {
+ // Permissive behavior: missing headers result in inclusion
+ if (headers == null) {
+ return true;
+ }
+
+ Header header = headers.lastHeader(headerName);
+
+ // Permissive behavior: header is missing, or its value is null or
zero-length. A zero-length value would
+ // otherwise decode to an empty string and be matched against the filter,
silently dropping such records.
+ if (header == null || header.value() == null || header.value().length ==
0) {
+ return true;
+ }
+
+ String headerValue = getDecodedHeaderValue(header.value());
+ // Permissive behavior: failed to decode header value
+ if (headerValue == null) {
+ return true;
+ }
+
+ return filterHandler.shouldInclude(headerValue);
+ }
+
+
+ /**
+ * Decode header bytes to string with caching.
+ * Returns null if decoding fails, so that undecodable header values fall
through to the permissive
+ * (include the record) path rather than matching on a replacement string.
+ */
+ @Nullable
+ private String getDecodedHeaderValue(byte[] headerBytes)
+ {
+ try {
+ // Do not cache oversized values so the cache cannot retain unbounded
heap for distinct large headers.
+ if (headerBytes.length > MAX_CACHEABLE_HEADER_BYTES) {
+ return decodeStrict(headerBytes);
+ }
+ ByteBuffer key = ByteBuffer.wrap(headerBytes);
+ return stringDecodingCache.get(key, k -> decodeStrict(headerBytes));
+ }
+ catch (Exception e) {
+ // Includes UncheckedIOException wrapping CharacterCodingException for
malformed/unmappable input.
+ log.warn(e, "Failed to decode header bytes, treating as null");
Review Comment:
[P2] Rate-limit malformed-header warnings
`getDecodedHeaderValue` is on the record-processing path, but this catch
logs the decoding exception at WARN, including a stack trace, for every
malformed or unmappable header value. A high-rate topic with bad bytes or an
encoding mismatch can therefore flood logs and spend substantial CPU and I/O on
logging, degrading ingestion and obscuring useful signals. Rate-limit or
deduplicate this warning, or lower the per-record path to debug while
preserving the permissive include behavior.
--
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]