From dca72ce65e09d74762ec45c69fb98ddc5d68af62 Mon Sep 17 00:00:00 2001 From: Marty Kulma <18468315+martykulma@users.noreply.github.com> Date: Wed, 30 Sep 2026 16:35:51 -0400 Subject: [PATCH 1/2] source: Kafka only comits offsets for a partition if they've changed --- src/storage/src/source/kafka.rs | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/src/storage/src/source/kafka.rs b/src/storage/src/source/kafka.rs index 05333c3a59898..2ad20ceea1d32 100644 --- a/src/storage/src/source/kafka.rs +++ b/src/storage/src/source/kafka.rs @@ -150,6 +150,9 @@ pub struct KafkaResumeUpperProcessor { config: RawSourceCreationConfig, topic_name: String, consumer: Arc>>, + /// The offsets this worker last committed upstream. Only a successful commit updates it, so a + /// failed commit is retried on the next resume upper. + committed_offsets: Vec<(PartitionId, MzOffset)>, } /// Computes whether this worker is responsible for consuming a partition. It assigns partitions to @@ -662,10 +665,11 @@ fn render_reader<'scope>( partition_capabilities, }; - let offset_committer = KafkaResumeUpperProcessor { + let mut offset_committer = KafkaResumeUpperProcessor { config: config.clone(), topic_name: topic.clone(), consumer, + committed_offsets: Vec::new(), }; // Seed the progress metrics from the resume uppers if we are snapshotting. @@ -1131,7 +1135,7 @@ fn render_reader<'scope>( impl KafkaResumeUpperProcessor { async fn process_frontier( - &self, + &mut self, uppers: &ResumeUppers, ) -> Result<(), anyhow::Error> { use rdkafka::consumer::CommitMode; @@ -1154,12 +1158,14 @@ impl KafkaResumeUpperProcessor { return Ok(()); }; let offsets: Vec<_> = self.responsible_offsets(frontier).collect(); - if !offsets.is_empty() { + // `uppers` changes whenever any export's upper moves, usually leaving this worker's + // offsets unchanged, and each commit is a synchronous round trip to the group coordinator. + if !offsets.is_empty() && offsets != self.committed_offsets { let mut tpl = TopicPartitionList::new(); - for (pid, offset) in offsets { + for (pid, offset) in &offsets { let offset_to_commit = Offset::Offset(offset.offset.try_into().expect("offset to be vald i64")); - tpl.add_partition_offset(&self.topic_name, pid, offset_to_commit) + tpl.add_partition_offset(&self.topic_name, *pid, offset_to_commit) .expect("offset known to be valid"); } let consumer = Arc::clone(&self.consumer); @@ -1168,6 +1174,7 @@ impl KafkaResumeUpperProcessor { move || consumer.commit(&tpl, CommitMode::Sync), ) .await?; + self.committed_offsets = offsets; } Ok(()) } From c2205069e7614e65f29f4171c70733107eecc9bd Mon Sep 17 00:00:00 2001 From: Marty Kulma <18468315+martykulma@users.noreply.github.com> Date: Fri, 2 Oct 2026 09:57:10 -0400 Subject: [PATCH 2/2] source: recommit unchanged Kafka offsets on a timer Brokers expire a standalone group's offset per partition offsets.retention.minutes after its last commit, so offsets the dedupe skips must still be recommitted periodically. The interval is the kafka_offset_commit_refresh_interval dyncfg, 10 minutes by default and 10 seconds in CI, with zero disabling the refresh. Co-Authored-By: Claude --- misc/python/materialize/mzcompose/__init__.py | 5 + .../materialize/parallel_workload/action.py | 1 + src/storage-types/src/dyncfgs.rs | 11 +++ src/storage/src/source/kafka.rs | 93 +++++++++++++------ .../mzcompose.py | 1 + 5 files changed, 85 insertions(+), 26 deletions(-) diff --git a/misc/python/materialize/mzcompose/__init__.py b/misc/python/materialize/mzcompose/__init__.py index c4ed34d524fca..f44ca7eae8887 100644 --- a/misc/python/materialize/mzcompose/__init__.py +++ b/misc/python/materialize/mzcompose/__init__.py @@ -422,6 +422,11 @@ def get_variable_system_parameters( "true" if force_source_table_syntax else "false", ["true", "false"] if force_source_table_syntax else ["false"], ), + # Low default so CI exercises the periodic recommit, which production + # only reaches after ten minutes. + VariableSystemParameter( + "kafka_offset_commit_refresh_interval", "10s", ["1s", "10s", "10min"] + ), VariableSystemParameter( "mysql_source_snapshot_parallelism", "true", ["true", "false"] ), diff --git a/misc/python/materialize/parallel_workload/action.py b/misc/python/materialize/parallel_workload/action.py index 6370d9a6ec38f..b430b33b3f03d 100644 --- a/misc/python/materialize/parallel_workload/action.py +++ b/misc/python/materialize/parallel_workload/action.py @@ -3400,6 +3400,7 @@ def __init__( "wallclock_lag_history_retention_interval", "wallclock_global_lag_histogram_retention_interval", "kafka_client_id_enrichment_rules", + "kafka_offset_commit_refresh_interval", "kafka_poll_max_wait", "kafka_default_aws_privatelink_endpoint_identification_algorithm", "kafka_buffered_event_resize_threshold_elements", diff --git a/src/storage-types/src/dyncfgs.rs b/src/storage-types/src/dyncfgs.rs index ebffb1602370a..126e810fab83b 100644 --- a/src/storage-types/src/dyncfgs.rs +++ b/src/storage-types/src/dyncfgs.rs @@ -114,6 +114,16 @@ pub const KAFKA_POLL_MAX_WAIT: Config = Config::new( ParameterScope::Replica, ); +/// How often Kafka sources recommit offsets that have not changed, so that brokers do not expire +/// them. Must stay well below the broker's `offsets.retention.minutes`. Zero disables it. +pub const KAFKA_OFFSET_COMMIT_REFRESH_INTERVAL: Config = Config::new( + "kafka_offset_commit_refresh_interval", + Duration::from_secs(10 * 60), + "How often Kafka sources recommit offsets that have not changed, which must stay well below \ + the broker's offsets.retention.minutes. Zero disables it.", + ParameterScope::Replica, +); + /// Whether to check the low watermark for Kafka sources and error if the start offset/resume /// upper has been compacted away. /// Environment-scoped because it decides whether a definite error is emitted. @@ -562,6 +572,7 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet { .add(&KAFKA_CLIENT_ID_ENRICHMENT_RULES) .add(&KAFKA_DEFAULT_AWS_PRIVATELINK_ENDPOINT_IDENTIFICATION_ALGORITHM) .add(&KAFKA_LOW_WATERMARK_CHECK) + .add(&KAFKA_OFFSET_COMMIT_REFRESH_INTERVAL) .add(&KAFKA_POLL_MAX_WAIT) .add(&KAFKA_RETRY_BACKOFF) .add(&KAFKA_RETRY_BACKOFF_MAX) diff --git a/src/storage/src/source/kafka.rs b/src/storage/src/source/kafka.rs index 2ad20ceea1d32..735ba93e705f9 100644 --- a/src/storage/src/source/kafka.rs +++ b/src/storage/src/source/kafka.rs @@ -700,17 +700,40 @@ fn render_reader<'scope>( } } + let refresh_interval = mz_storage_types::dyncfgs::KAFKA_OFFSET_COMMIT_REFRESH_INTERVAL + .get(config.config.config_set()); let resume_uppers_process_loop = async move { tokio::pin!(resume_uppers); - while let Some(uppers) = resume_uppers.next().await { - if let Err(e) = offset_committer.process_frontier(&uppers).await { - offset_commit_metrics.offset_commit_failures.inc(); - tracing::warn!( - %e, - "timely-{worker_id} source({source_id}) failed to commit offsets: {uppers}", - worker_id = config.worker_id, - source_id = config.id, - ); + // Zero disables the refresh. `interval` panics on a zero period, so the tick arm + // below is what honors it. + let mut refresh = + tokio::time::interval(refresh_interval.max(Duration::from_millis(1))); + // A commit that outlasts a period must not be followed by a burst of recommits. + refresh.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + let report = |e: anyhow::Error, what: &dyn std::fmt::Display| { + offset_commit_metrics.offset_commit_failures.inc(); + tracing::warn!( + %e, + "timely-{worker_id} source({source_id}) failed to commit offsets: {what}", + worker_id = config.worker_id, + source_id = config.id, + ); + }; + loop { + tokio::select! { + uppers = resume_uppers.next() => match uppers { + Some(uppers) => { + if let Err(e) = offset_committer.process_frontier(&uppers).await { + report(e, &uppers); + } + } + None => break, + }, + _ = refresh.tick(), if !refresh_interval.is_zero() => { + if let Err(e) = offset_committer.refresh().await { + report(e, &"refresh"); + } + } } } // During dataflow shutdown this loop can end due to the general chaos caused by @@ -1134,12 +1157,12 @@ fn render_reader<'scope>( } impl KafkaResumeUpperProcessor { + /// Reports each export's `offset_committed` and commits the source upper's offsets upstream + /// when they differ from the last successful commit. async fn process_frontier( &mut self, uppers: &ResumeUppers, ) -> Result<(), anyhow::Error> { - use rdkafka::consumer::CommitMode; - for (id, frontier) in &uppers.exports { if let Some(stat) = self.config.statistics.get(id) { // Note that we do not subtract 1 from the frontier. Imagine @@ -1160,22 +1183,40 @@ impl KafkaResumeUpperProcessor { let offsets: Vec<_> = self.responsible_offsets(frontier).collect(); // `uppers` changes whenever any export's upper moves, usually leaving this worker's // offsets unchanged, and each commit is a synchronous round trip to the group coordinator. - if !offsets.is_empty() && offsets != self.committed_offsets { - let mut tpl = TopicPartitionList::new(); - for (pid, offset) in &offsets { - let offset_to_commit = - Offset::Offset(offset.offset.try_into().expect("offset to be vald i64")); - tpl.add_partition_offset(&self.topic_name, *pid, offset_to_commit) - .expect("offset known to be valid"); - } - let consumer = Arc::clone(&self.consumer); - mz_ore::task::spawn_blocking( - || format!("source({}) kafka offset commit", self.config.id), - move || consumer.commit(&tpl, CommitMode::Sync), - ) - .await?; - self.committed_offsets = offsets; + if offsets.is_empty() || offsets == self.committed_offsets { + return Ok(()); + } + self.commit(&offsets).await?; + self.committed_offsets = offsets; + Ok(()) + } + + /// Recommits the offsets of the last successful commit. The consumer only ever `assign`s, so + /// the broker treats its group as standalone and expires a partition's offset + /// `offsets.retention.minutes` after that partition's last commit, current or not. + async fn refresh(&self) -> Result<(), anyhow::Error> { + if self.committed_offsets.is_empty() { + return Ok(()); } + self.commit(&self.committed_offsets).await + } + + async fn commit(&self, offsets: &[(PartitionId, MzOffset)]) -> Result<(), anyhow::Error> { + use rdkafka::consumer::CommitMode; + + let mut tpl = TopicPartitionList::new(); + for (pid, offset) in offsets { + let offset_to_commit = + Offset::Offset(offset.offset.try_into().expect("offset to be vald i64")); + tpl.add_partition_offset(&self.topic_name, *pid, offset_to_commit) + .expect("offset known to be valid"); + } + let consumer = Arc::clone(&self.consumer); + mz_ore::task::spawn_blocking( + || format!("source({}) kafka offset commit", self.config.id), + move || consumer.commit(&tpl, CommitMode::Sync), + ) + .await?; Ok(()) } diff --git a/test/launchdarkly-flag-consistency/mzcompose.py b/test/launchdarkly-flag-consistency/mzcompose.py index 059dbf0650eb1..af4452ccf1f66 100644 --- a/test/launchdarkly-flag-consistency/mzcompose.py +++ b/test/launchdarkly-flag-consistency/mzcompose.py @@ -306,6 +306,7 @@ hydration_history_retention_period kafka_buffered_event_resize_threshold_elements kafka_default_aws_privatelink_endpoint_identification_algorithm + kafka_offset_commit_refresh_interval kafka_poll_max_wait kafka_reconnect_backoff kafka_reconnect_backoff_max