Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions misc/python/materialize/mzcompose/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
),
Expand Down
1 change: 1 addition & 0 deletions misc/python/materialize/parallel_workload/action.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
11 changes: 11 additions & 0 deletions src/storage-types/src/dyncfgs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,16 @@ pub const KAFKA_POLL_MAX_WAIT: Config<Duration> = 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<Duration> = 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.
Expand Down Expand Up @@ -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)
Expand Down
102 changes: 75 additions & 27 deletions src/storage/src/source/kafka.rs
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,9 @@ pub struct KafkaResumeUpperProcessor {
config: RawSourceCreationConfig,
topic_name: String,
consumer: Arc<BaseConsumer<TunnelingClientContext<GlueConsumerContext>>>,
/// 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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -696,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)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

if the value was set to 0, and this clamps it to 1ms, is that sustainable? I suppose it would just happen every iteration of the loop, which is okay?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1ms would just update as fast as the calls would allow it (because the missed tick behavior is delay).

If someone sets it to 0 to disable, then the branch never fires because the precondition would return false:

_ = refresh.tick(), if !refresh_interval.is_zero() => {
if let Err(e) = offset_committer.refresh().await {
report(e, &"refresh");
}

https://docs.rs/tokio/latest/tokio/macro.select.html

Evaluate all provided expressions. If the precondition returns false, disable the branch for the remainder of the current call to select!. Re-entering select! due to a loop clears the “disabled” state.

// 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
Expand Down Expand Up @@ -1130,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(
&self,
&mut self,
uppers: &ResumeUppers<KafkaTimestamp>,
) -> 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
Expand All @@ -1154,21 +1181,42 @@ impl KafkaResumeUpperProcessor {
return Ok(());
};
let offsets: Vec<_> = self.responsible_offsets(frontier).collect();
if !offsets.is_empty() {
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?;
// `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 {
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(())
}

Expand Down
1 change: 1 addition & 0 deletions test/launchdarkly-flag-consistency/mzcompose.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading