Skip to content

Commit f54cbcc

Browse files
committed
[adapters] Kafka header filters.
New Kafka input connector config setting to filter messages based on Kafka message headers. Supports primitive checks against regular expressions and arbitrary Boolean combinations. Signed-off-by: Leonid Ryzhyk <ryzhyk@gmail.com>
1 parent a3370ff commit f54cbcc

9 files changed

Lines changed: 919 additions & 8 deletions

File tree

Cargo.lock

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

crates/adapters/src/transport/kafka/ft/input.rs

Lines changed: 55 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,9 @@ use feldera_adapterlib::transport::{
2323
use feldera_sqllib::{ByteArray, SqlString, Timestamp, Variant};
2424
use feldera_types::config::FtModel;
2525
use feldera_types::program_schema::Relation;
26-
use feldera_types::transport::kafka::{KafkaInputConfig, KafkaStartFromConfig};
26+
use feldera_types::transport::kafka::{
27+
CompiledHeaderFilter, KafkaInputConfig, KafkaStartFromConfig,
28+
};
2729
use itertools::Itertools;
2830
use rdkafka::client::OAuthToken;
2931
use rdkafka::config::RDKafkaLogLevel;
@@ -66,6 +68,10 @@ const METADATA_TIMEOUT: Duration = Duration::from_secs(10);
6668
// to the worker thread.
6769
const ERROR_BUFFER_SIZE: usize = 1000;
6870

71+
/// A Kafka message's headers borrowed as `(key, value)` pairs for header
72+
/// filtering. A `None` value denotes a header present with a null value.
73+
type HeaderPairs<'a> = SmallVec<[(&'a str, Option<&'a [u8]>); 8]>;
74+
6975
pub struct KafkaFtInputEndpoint {
7076
config: Arc<KafkaInputConfig>,
7177
}
@@ -380,6 +386,16 @@ impl KafkaFtInputReaderInner {
380386
})
381387
.collect::<Vec<_>>();
382388

389+
// Compile the header filter once and share it across partition
390+
// receivers. A malformed filter is already rejected by
391+
// `KafkaInputConfig::validate`, so this normally cannot fail.
392+
let header_filter = config
393+
.header_filter
394+
.as_ref()
395+
.map(|filter| filter.compile())
396+
.transpose()?
397+
.map(Arc::new);
398+
383399
// Split every partition away as its own separate queue.
384400
let mut receivers = BTreeMap::new();
385401

@@ -405,6 +421,7 @@ impl KafkaFtInputReaderInner {
405421
queue,
406422
next_offset,
407423
&config,
424+
header_filter.clone(),
408425
unparker,
409426
));
410427
receivers.insert(partition, receiver.clone());
@@ -999,6 +1016,10 @@ struct PartitionReceiver {
9991016
config: KafkaInputConfig,
10001017
metadata_requested: bool,
10011018

1019+
/// Compiled header filter, shared across partitions. Messages that do not
1020+
/// match are dropped before parsing. `None` admits all messages.
1021+
header_filter: Option<Arc<CompiledHeaderFilter>>,
1022+
10021023
/// The maximum message offset that we want to receive, used as follows:
10031024
///
10041025
/// - `i64::MIN`, the initial value, disables receiving messages entirely.
@@ -1041,6 +1062,7 @@ impl PartitionReceiver {
10411062
queue: PartitionQueue<KafkaFtInputContext>,
10421063
next_offset: i64,
10431064
config: &KafkaInputConfig,
1065+
header_filter: Option<Arc<CompiledHeaderFilter>>,
10441066
unparker: Unparker,
10451067
) -> Self {
10461068
let metadata_requested = config.metadata_requested();
@@ -1057,6 +1079,7 @@ impl PartitionReceiver {
10571079
fatal_error: AtomicBool::new(false),
10581080
config: config.clone(),
10591081
metadata_requested,
1082+
header_filter,
10601083
unparker,
10611084
}
10621085
}
@@ -1112,6 +1135,27 @@ impl PartitionReceiver {
11121135
self.next_offset.load(Ordering::Relaxed)
11131136
}
11141137

1138+
/// Returns `true` if `message` satisfies the configured header filter, or if
1139+
/// no filter is configured. Messages for which this returns `false` are
1140+
/// dropped before parsing.
1141+
fn header_filter_admits(&self, message: &BorrowedMessage<'_>) -> bool {
1142+
let Some(filter) = &self.header_filter else {
1143+
return true;
1144+
};
1145+
1146+
// Borrow the message's headers as (key, value) pairs; nothing is copied.
1147+
let headers: HeaderPairs = match message.headers() {
1148+
Some(headers) => (0..headers.count())
1149+
.map(|i| {
1150+
let header = headers.get(i);
1151+
(header.key, header.value)
1152+
})
1153+
.collect(),
1154+
None => SmallVec::new(),
1155+
};
1156+
filter.matches(&headers)
1157+
}
1158+
11151159
/// Create record metadata from Kafka message containing only properties specified in the connector config.
11161160
fn create_metadata(&self, message: &BorrowedMessage<'_>) -> Option<ConnectorMetadata> {
11171161
if !self.metadata_requested {
@@ -1189,6 +1233,16 @@ impl PartitionReceiver {
11891233
let next_offset = self.next_offset();
11901234
if offset >= next_offset {
11911235
self.next_offset.store(offset + 1, Ordering::Relaxed);
1236+
1237+
// Drop messages rejected by the header filter before parsing,
1238+
// so they never enter the offset ranges or the step hash.
1239+
// This runs after advancing `next_offset` (the consume
1240+
// position still moves forward) and is re-applied identically
1241+
// on replay, keeping ranges and hashes deterministic.
1242+
if !self.header_filter_admits(&message) {
1243+
return;
1244+
}
1245+
11921246
let timestamp = message.timestamp().to_millis().unwrap_or(i64::MIN);
11931247
let payload = message.payload().unwrap_or(&[]);
11941248
let metadata = self.create_metadata(&message);

crates/adapters/src/transport/kafka/ft/test.rs

Lines changed: 179 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -136,14 +136,37 @@ fn create_reader(
136136
DummyInputReceiver,
137137
Box<dyn InputReader>,
138138
) {
139-
let config = serde_json::from_value(json!({
140-
"name": "kafka_input",
141-
"config": {
139+
create_reader_config(topic, resume_info, synchronize_partitions, None)
140+
}
141+
142+
/// Like [`create_reader`], but with an optional `header_filter` added to the
143+
/// connector configuration.
144+
fn create_reader_config(
145+
topic: &str,
146+
resume_info: Option<JsonValue>,
147+
synchronize_partitions: bool,
148+
header_filter: Option<JsonValue>,
149+
) -> (
150+
Box<dyn TransportInputEndpoint>,
151+
DummyInputReceiver,
152+
Box<dyn InputReader>,
153+
) {
154+
let mut inner_config = json!({
142155
"topic": topic,
143156
"log_level": "debug",
144-
"start_from": "earliest",
145-
"synchronize_partitions": synchronize_partitions,
146-
}
157+
"start_from": "earliest",
158+
"synchronize_partitions": synchronize_partitions,
159+
});
160+
if let Some(header_filter) = header_filter {
161+
inner_config
162+
.as_object_mut()
163+
.unwrap()
164+
.insert("header_filter".to_string(), header_filter);
165+
}
166+
167+
let config = serde_json::from_value(json!({
168+
"name": "kafka_input",
169+
"config": inner_config,
147170
}))
148171
.unwrap();
149172

@@ -165,6 +188,23 @@ fn create_reader(
165188
(endpoint, receiver, reader)
166189
}
167190

191+
/// Drain and return the payloads flushed to `receiver` so far, in order.
192+
fn take_flushed(receiver: &DummyInputReceiver) -> Vec<String> {
193+
mem::take(&mut *receiver.inner.flushed.lock().unwrap())
194+
}
195+
196+
/// Build the `headers` argument for [`TestProducer::send_message`] from
197+
/// `(key, value)` pairs, where a `None` value denotes a header with a null
198+
/// value.
199+
fn headers(pairs: &[(&str, Option<&[u8]>)]) -> Option<BTreeMap<String, Option<Vec<u8>>>> {
200+
Some(
201+
pairs
202+
.iter()
203+
.map(|(key, value)| (key.to_string(), value.map(|v| v.to_vec())))
204+
.collect(),
205+
)
206+
}
207+
168208
#[test]
169209
fn single_input() {
170210
test_input("single_input_ft", &[10]);
@@ -469,6 +509,139 @@ fn test_input(topic: &str, batch_sizes: &[u32]) {
469509
}
470510
}
471511

512+
/// A configured header filter drops non-matching messages before parsing, and
513+
/// the drop decision is deterministic across replay: a fresh reader replaying
514+
/// the checkpointed offset range reproduces exactly the admitted records.
515+
///
516+
/// To confirm this test catches a regression, make `header_filter_admits`
517+
/// always return `true`; the connector then buffers all 6 messages instead of
518+
/// 3 and both the live and the replay expectations fail.
519+
#[test]
520+
fn test_input_header_filter() {
521+
init_test_logger();
522+
523+
let topic = "kafka_header_filter_ft";
524+
let _kafka_resources = KafkaResources::create_topics(&[(topic, 1)]);
525+
526+
// Admit only messages whose `keep` header equals `yes`.
527+
let filter = json!({"header": {"name": "keep", "pattern": "yes"}});
528+
529+
let (endpoint, receiver, reader) =
530+
create_reader_config(topic, None, false, Some(filter.clone()));
531+
reader.extend();
532+
533+
// Offsets 0..6. `keep=yes` at 0, 2, 4 are admitted; 1, 3 (`keep=no`) and 5
534+
// (no header) are dropped. The trailing drop at offset 5 lies beyond the
535+
// recorded range and must never be replayed.
536+
let producer = TestProducer::new();
537+
producer.send_message(b"m0", topic, headers(&[("keep", Some(b"yes"))]));
538+
producer.send_message(b"m1", topic, headers(&[("keep", Some(b"no"))]));
539+
producer.send_message(b"m2", topic, headers(&[("keep", Some(b"yes"))]));
540+
producer.send_message(b"m3", topic, headers(&[("keep", Some(b"no"))]));
541+
producer.send_message(b"m4", topic, headers(&[("keep", Some(b"yes"))]));
542+
producer.send_message(b"m5", topic, None);
543+
544+
// Only the 3 admitted messages are buffered.
545+
receiver.expect_buffering(3);
546+
reader.queue(false);
547+
548+
// The offset range spans the first through the last admitted offset (0..5);
549+
// dropped offsets inside and beyond the range do not change it.
550+
let metadata = Metadata {
551+
offsets: vec![0..5],
552+
};
553+
receiver.expect(vec![ConsumerCall::Extended {
554+
num_records: 3,
555+
metadata: serde_json::to_value(&metadata).unwrap(),
556+
}]);
557+
assert_eq!(take_flushed(&receiver), vec!["m0", "m2", "m4"]);
558+
559+
drop(endpoint);
560+
drop(receiver);
561+
drop(reader);
562+
563+
// Replay the recorded range with the same filter. The filter is re-applied
564+
// during replay, so exactly the admitted records reappear, in order.
565+
let (_endpoint, receiver, reader) = create_reader_config(topic, None, false, Some(filter));
566+
receiver.inner.drop_buffered.store(true, Ordering::Release);
567+
reader.replay(serde_json::to_value(&metadata).unwrap(), RmpValue::Nil);
568+
receiver.expect(vec![ConsumerCall::Replayed { num_records: 3 }]);
569+
assert_eq!(take_flushed(&receiver), vec!["m0", "m2", "m4"]);
570+
}
571+
572+
/// A boolean header filter (`and`/`or`/`not`) drops the right messages in the
573+
/// real connector, exercising leading, middle, and trailing drops, and its
574+
/// decisions are re-applied identically on replay. The filter reads headers
575+
/// directly, so it works even though `include_headers` is not set.
576+
#[test]
577+
fn test_input_header_filter_boolean() {
578+
init_test_logger();
579+
580+
let topic = "kafka_header_filter_boolean_ft";
581+
let _kafka_resources = KafkaResources::create_topics(&[(topic, 1)]);
582+
583+
// Admit iff (env is prod or staging) and not (skip == drop).
584+
let filter = json!({
585+
"and": [
586+
{"or": [
587+
{"header": {"name": "env", "pattern": "prod"}},
588+
{"header": {"name": "env", "pattern": "staging"}}
589+
]},
590+
{"not": {"header": {"name": "skip", "pattern": "drop"}}}
591+
]
592+
});
593+
594+
let (endpoint, receiver, reader) =
595+
create_reader_config(topic, None, false, Some(filter.clone()));
596+
reader.extend();
597+
598+
let producer = TestProducer::new();
599+
// 0: env=dev -> drop (fails `or`); leading drop
600+
// 1: env=prod -> admit
601+
// 2: env=staging, skip=drop -> drop (fails `not`); middle drop
602+
// 3: env=staging -> admit
603+
// 4: env=prod, skip=drop -> drop (fails `not`); trailing drop
604+
producer.send_message(b"m0", topic, headers(&[("env", Some(b"dev"))]));
605+
producer.send_message(b"m1", topic, headers(&[("env", Some(b"prod"))]));
606+
producer.send_message(
607+
b"m2",
608+
topic,
609+
headers(&[("env", Some(b"staging")), ("skip", Some(b"drop"))]),
610+
);
611+
producer.send_message(b"m3", topic, headers(&[("env", Some(b"staging"))]));
612+
producer.send_message(
613+
b"m4",
614+
topic,
615+
headers(&[("env", Some(b"prod")), ("skip", Some(b"drop"))]),
616+
);
617+
618+
// Admitted: offsets 1 and 3.
619+
receiver.expect_buffering(2);
620+
reader.queue(false);
621+
622+
// The range starts at the first admitted offset (1) and ends after the last
623+
// (3), skipping the leading drop at offset 0.
624+
let metadata = Metadata {
625+
offsets: vec![1..4],
626+
};
627+
receiver.expect(vec![ConsumerCall::Extended {
628+
num_records: 2,
629+
metadata: serde_json::to_value(&metadata).unwrap(),
630+
}]);
631+
assert_eq!(take_flushed(&receiver), vec!["m1", "m3"]);
632+
633+
drop(endpoint);
634+
drop(receiver);
635+
drop(reader);
636+
637+
// Replaying the recorded range reproduces the admitted records.
638+
let (_endpoint, receiver, reader) = create_reader_config(topic, None, false, Some(filter));
639+
receiver.inner.drop_buffered.store(true, Ordering::Release);
640+
reader.replay(serde_json::to_value(&metadata).unwrap(), RmpValue::Nil);
641+
receiver.expect(vec![ConsumerCall::Replayed { num_records: 2 }]);
642+
assert_eq!(take_flushed(&receiver), vec!["m1", "m3"]);
643+
}
644+
472645
#[derive(Debug, PartialEq)]
473646
enum ConsumerCall {
474647
ParseErrors,

crates/feldera-types/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,3 +40,4 @@ clap = { workspace = true }
4040
[dev-dependencies]
4141
csv = { workspace = true }
4242
tempfile = { workspace = true }
43+
serde_yaml = { workspace = true }

0 commit comments

Comments
 (0)