1use std::collections::{HashMap, HashSet};
16use std::sync::{Arc, LazyLock, Weak};
17use std::time::Duration;
18
19use anyhow::{Context, anyhow};
20use async_trait::async_trait;
21use moka::future::Cache as MokaCache;
22use moka::ops::compute::Op;
23use rdkafka::admin::{AdminClient, AdminOptions};
24use rdkafka::consumer::BaseConsumer;
25#[cfg(not(madsim))]
26use rdkafka::consumer::Consumer;
27use rdkafka::error::{KafkaError, KafkaResult};
28use rdkafka::types::RDKafkaErrorCode;
29use rdkafka::{ClientConfig, Offset, TopicPartitionList};
30use risingwave_common::bail;
31use risingwave_common::id::FragmentId;
32use risingwave_common::metrics::LabelGuardedIntGauge;
33use thiserror_ext::AsReport;
34
35use crate::connector_common::read_kafka_log_level;
36use crate::error::{ConnectorError, ConnectorResult};
37use crate::source::SourceEnumeratorContextRef;
38use crate::source::base::SplitEnumerator;
39use crate::source::kafka::split::KafkaSplit;
40use crate::source::kafka::{
41 KAFKA_ISOLATION_LEVEL, KafkaConnectionProps, KafkaContextCommon, KafkaProperties,
42 RwConsumerContext,
43};
44
45type KafkaConsumer = BaseConsumer<RwConsumerContext>;
46type KafkaAdmin = AdminClient<RwConsumerContext>;
47
48pub static SHARED_KAFKA_CONSUMER: LazyLock<MokaCache<KafkaConnectionProps, Weak<KafkaConsumer>>> =
50 LazyLock::new(|| moka::future::Cache::builder().build());
51pub static SHARED_KAFKA_ADMIN: LazyLock<MokaCache<KafkaConnectionProps, Arc<KafkaAdmin>>> =
53 LazyLock::new(|| {
54 moka::future::Cache::builder()
55 .time_to_idle(Duration::from_secs(5 * 60))
56 .build()
57 });
58
59#[derive(Debug, Copy, Clone, Eq, PartialEq)]
60pub enum KafkaEnumeratorOffset {
61 Earliest,
62 Latest,
63 Timestamp(i64),
64 None,
65}
66
67pub struct KafkaSplitEnumerator {
68 context: SourceEnumeratorContextRef,
69 broker_address: String,
70 topic: String,
71 client: Arc<KafkaConsumer>,
72 start_offset: KafkaEnumeratorOffset,
73 resolved_start_offsets: HashMap<i32, Option<i64>>,
80
81 stop_offset: KafkaEnumeratorOffset,
83
84 sync_call_timeout: Duration,
85 high_watermark_metrics: HashMap<i32, LabelGuardedIntGauge>,
86
87 properties: KafkaProperties,
88 config: rdkafka::ClientConfig,
89}
90
91impl KafkaSplitEnumerator {
92 fn report_consumer_group_delete_failure(&self, group_id: &str) {
93 let source_id = self.context.info.source_id.to_string();
94 self.context
95 .metrics
96 .kafka_consumer_group_delete_failure_count
97 .with_label_values(&[&source_id, group_id])
98 .inc();
99 }
100
101 async fn drop_consumer_groups(&self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
102 let admin = Box::pin(SHARED_KAFKA_ADMIN.try_get_with_by_ref(
103 &self.properties.connection,
104 async {
105 tracing::info!("build new kafka admin for {}", self.broker_address);
106 Ok(Arc::new(
107 build_kafka_admin(&self.config, &self.properties).await?,
108 ))
109 },
110 ))
111 .await?;
112
113 let group_ids = fragment_ids
114 .iter()
115 .map(|fragment_id| self.properties.group_id(*fragment_id))
116 .collect::<Vec<_>>();
117 let group_id_refs = group_ids
118 .iter()
119 .map(|group_id| group_id.as_str())
120 .collect::<Vec<_>>();
121
122 let res = match admin
123 .delete_groups(&group_id_refs, &AdminOptions::default())
124 .await
125 {
126 Ok(res) => res,
127 Err(err) => {
128 for group_id in &group_ids {
129 self.report_consumer_group_delete_failure(group_id);
130 }
131 tracing::warn!(
132 error = %err.as_report(),
133 topic = self.topic,
134 ?group_ids,
135 "failed to delete Kafka consumer groups"
136 );
137 return Err(err.into());
138 }
139 };
140
141 let mut failure_count = 0;
142 for result in &res {
143 if let Err((group_id, error_code)) = result {
144 failure_count += 1;
145 self.report_consumer_group_delete_failure(group_id);
146 tracing::warn!(
147 topic = self.topic,
148 group_id,
149 error = %error_code.as_report(),
150 "failed to delete Kafka consumer group"
151 );
152 }
153 }
154 tracing::debug!(
155 topic = self.topic,
156 ?fragment_ids,
157 ?res,
158 failure_count,
159 "delete groups result"
160 );
161 Ok(())
162 }
163}
164
165#[async_trait]
166impl SplitEnumerator for KafkaSplitEnumerator {
167 type Properties = KafkaProperties;
168 type Split = KafkaSplit;
169
170 async fn new(
171 properties: KafkaProperties,
172 context: SourceEnumeratorContextRef,
173 ) -> ConnectorResult<KafkaSplitEnumerator> {
174 let mut config = rdkafka::ClientConfig::new();
175 let common_props = &properties.common;
176
177 let broker_address = properties.connection.brokers.clone();
178 let topic = common_props.topic.clone();
179 config.set("bootstrap.servers", &broker_address);
180 config.set("isolation.level", KAFKA_ISOLATION_LEVEL);
181 if let Some(log_level) = read_kafka_log_level() {
182 config.set_log_level(log_level);
183 }
184 properties.connection.set_security_properties(&mut config);
185 properties.set_client(&mut config);
186 config.set("statistics.interval.ms", "0");
190 let mut scan_start_offset = match properties
191 .scan_startup_mode
192 .as_ref()
193 .map(|s| s.to_lowercase())
194 .as_deref()
195 {
196 Some("earliest") => KafkaEnumeratorOffset::Earliest,
197 Some("latest") => KafkaEnumeratorOffset::Latest,
198 None => KafkaEnumeratorOffset::Earliest,
199 _ => bail!(
200 "properties `scan_startup_mode` only supports earliest and latest or leaving it empty"
201 ),
202 };
203
204 if let Some(s) = &properties.time_offset {
205 let time_offset = s.parse::<i64>().map_err(|e| anyhow!(e))?;
206 scan_start_offset = KafkaEnumeratorOffset::Timestamp(time_offset)
207 }
208
209 let mut client: Option<Arc<KafkaConsumer>> = None;
210 SHARED_KAFKA_CONSUMER
211 .entry_by_ref(&properties.connection)
212 .and_try_compute_with::<_, _, ConnectorError>(|maybe_entry| async {
213 if let Some(entry) = maybe_entry {
214 let entry_value = entry.into_value();
215 if let Some(client_) = entry_value.upgrade() {
216 tracing::info!("reuse existing kafka client for {}", broker_address);
218 client = Some(client_);
219 return Ok(Op::Nop);
220 }
221 }
222 tracing::info!("build new kafka client for {}", broker_address);
223 client = Some(build_kafka_client(&config, &properties).await?);
224 Ok(Op::Put(Arc::downgrade(client.as_ref().unwrap())))
225 })
226 .await?;
227
228 Ok(Self {
229 context,
230 broker_address,
231 topic,
232 client: client.unwrap(),
233 start_offset: scan_start_offset,
234 resolved_start_offsets: HashMap::new(),
235 stop_offset: KafkaEnumeratorOffset::None,
236 sync_call_timeout: properties.common.sync_call_timeout,
237 high_watermark_metrics: HashMap::new(),
238 properties,
239 config,
240 })
241 }
242
243 async fn list_splits(&mut self) -> ConnectorResult<Vec<KafkaSplit>> {
244 if let Some(Err(poll_err)) = {
253 #[cfg(not(madsim))]
254 {
255 self.client.poll(Duration::ZERO)
256 }
257 #[cfg(madsim)]
258 {
259 self.client.poll(Duration::ZERO).await
260 }
261 } && !is_expected_no_group_poll_error(&poll_err)
262 {
263 tracing::warn!(
264 error = %poll_err.as_report(),
265 topic = self.topic,
266 broker_address = self.broker_address,
267 "failed to poll kafka client");
268 }
269
270 let topic_partitions = self.fetch_topic_partition().await.with_context(|| {
271 format!(
272 "failed to fetch metadata from kafka ({})",
273 self.broker_address
274 )
275 })?;
276
277 let watermarks = self.get_watermarks(topic_partitions.as_ref()).await?;
278 let mut start_offsets = self
279 .fetch_start_offset(topic_partitions.as_ref(), &watermarks)
280 .await?;
281
282 let mut stop_offsets = self
283 .fetch_stop_offset(topic_partitions.as_ref(), &watermarks)
284 .await?;
285
286 let ret: Vec<_> = topic_partitions
287 .into_iter()
288 .map(|partition| KafkaSplit {
289 topic: self.topic.clone(),
290 partition,
291 start_offset: start_offsets.remove(&partition).unwrap(),
292 stop_offset: stop_offsets.remove(&partition).unwrap(),
293 })
294 .collect();
295
296 Ok(ret)
297 }
298
299 async fn on_drop_fragments(&mut self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
300 self.drop_consumer_groups(fragment_ids).await
301 }
302
303 async fn on_finish_backfill(&mut self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
304 self.drop_consumer_groups(fragment_ids).await
305 }
306}
307
308fn is_expected_no_group_poll_error(error: &KafkaError) -> bool {
309 matches!(
310 error,
311 KafkaError::MessageConsumption(RDKafkaErrorCode::UnknownGroup)
312 )
313}
314
315async fn build_kafka_client(
316 config: &ClientConfig,
317 properties: &KafkaProperties,
318) -> ConnectorResult<Arc<KafkaConsumer>> {
319 let ctx_common = KafkaContextCommon::new(
320 properties.privatelink_common.broker_rewrite_map.clone(),
321 None,
322 None,
323 properties.aws_auth_props.clone(),
324 properties.connection.is_aws_msk_iam(),
325 )
326 .await?;
327 let client_ctx = RwConsumerContext::new(ctx_common);
328 let client: KafkaConsumer = config.create_with_context(client_ctx).await?;
329
330 if properties.connection.is_aws_msk_iam() {
336 #[cfg(not(madsim))]
337 client.poll(Duration::from_secs(10)); #[cfg(madsim)]
339 client.poll(Duration::from_secs(10)).await;
340 }
341 Ok(Arc::new(client))
342}
343async fn build_kafka_admin(
344 config: &ClientConfig,
345 properties: &KafkaProperties,
346) -> ConnectorResult<KafkaAdmin> {
347 let ctx_common = KafkaContextCommon::new(
348 properties.privatelink_common.broker_rewrite_map.clone(),
349 None,
350 None,
351 properties.aws_auth_props.clone(),
352 properties.connection.is_aws_msk_iam(),
353 )
354 .await?;
355 let client_ctx = RwConsumerContext::new(ctx_common);
356 let client: KafkaAdmin = config.create_with_context(client_ctx).await?;
357 Ok(client)
359}
360
361impl KafkaSplitEnumerator {
362 async fn get_watermarks(
363 &mut self,
364 partitions: &[i32],
365 ) -> KafkaResult<HashMap<i32, (i64, i64)>> {
366 let mut map = HashMap::new();
367 for partition in partitions {
368 let (low, high) = self
369 .client
370 .fetch_watermarks(self.topic.as_str(), *partition, self.sync_call_timeout)
371 .await?;
372 self.report_high_watermark(*partition, high);
373 map.insert(*partition, (low, high));
374 }
375 tracing::debug!("fetch kafka watermarks: {map:?}");
376 Ok(map)
377 }
378
379 pub async fn list_splits_batch(
380 &mut self,
381 expect_start_timestamp_millis: Option<i64>,
382 expect_stop_timestamp_millis: Option<i64>,
383 ) -> ConnectorResult<Vec<KafkaSplit>> {
384 let topic_partitions = self.fetch_topic_partition().await.with_context(|| {
385 format!(
386 "failed to fetch metadata from kafka ({})",
387 self.broker_address
388 )
389 })?;
390
391 let mut watermarks = self.get_watermarks(topic_partitions.as_ref()).await?;
394
395 let mut expect_start_offset = if let Some(ts) = expect_start_timestamp_millis {
400 Some(
401 self.fetch_offset_for_time(topic_partitions.as_ref(), ts, &watermarks)
402 .await?,
403 )
404 } else {
405 None
406 };
407
408 let mut expect_stop_offset = if let Some(ts) = expect_stop_timestamp_millis {
409 Some(
410 self.fetch_offset_for_time(topic_partitions.as_ref(), ts, &watermarks)
411 .await?,
412 )
413 } else {
414 None
415 };
416
417 Ok(topic_partitions
418 .iter()
419 .map(|partition| {
420 let (low, high) = watermarks.remove(partition).unwrap();
421 let start_offset = {
422 let earliest_offset = low - 1;
423 let start = expect_start_offset
424 .as_mut()
425 .map(|m| m.remove(partition).flatten().unwrap_or(earliest_offset))
426 .unwrap_or(earliest_offset);
427 i64::max(start, earliest_offset)
428 };
429 let stop_offset = {
430 let stop = expect_stop_offset
431 .as_mut()
432 .map(|m| m.remove(partition).unwrap_or(Some(high)))
433 .unwrap_or(Some(high))
434 .unwrap_or(high);
435 i64::min(stop, high)
436 };
437
438 if start_offset > stop_offset {
439 tracing::warn!(
440 "Skipping topic {} partition {}: requested start offset {} is greater than stop offset {}",
441 self.topic,
442 partition,
443 start_offset,
444 stop_offset
445 );
446 }
447 KafkaSplit {
448 topic: self.topic.clone(),
449 partition: *partition,
450 start_offset: Some(start_offset),
451 stop_offset: Some(stop_offset),
452 }
453 })
454 .collect::<Vec<KafkaSplit>>())
455 }
456
457 async fn fetch_stop_offset(
458 &self,
459 partitions: &[i32],
460 watermarks: &HashMap<i32, (i64, i64)>,
461 ) -> KafkaResult<HashMap<i32, Option<i64>>> {
462 match self.stop_offset {
463 KafkaEnumeratorOffset::Earliest => unreachable!(),
464 KafkaEnumeratorOffset::Latest => {
465 let mut map = HashMap::new();
466 for partition in partitions {
467 let (_, high_watermark) = watermarks.get(partition).unwrap();
468 map.insert(*partition, Some(*high_watermark));
469 }
470 Ok(map)
471 }
472 KafkaEnumeratorOffset::Timestamp(time) => {
473 self.fetch_offset_for_time(partitions, time, watermarks)
474 .await
475 }
476 KafkaEnumeratorOffset::None => partitions
477 .iter()
478 .map(|partition| Ok((*partition, None)))
479 .collect(),
480 }
481 }
482
483 async fn fetch_start_offset(
484 &mut self,
485 partitions: &[i32],
486 watermarks: &HashMap<i32, (i64, i64)>,
487 ) -> KafkaResult<HashMap<i32, Option<i64>>> {
488 match self.start_offset {
489 KafkaEnumeratorOffset::Earliest | KafkaEnumeratorOffset::Latest => {
490 let mut map = HashMap::new();
491 for partition in partitions {
492 let (low_watermark, high_watermark) = watermarks.get(partition).unwrap();
493 let offset = match self.start_offset {
494 KafkaEnumeratorOffset::Earliest => low_watermark - 1,
495 KafkaEnumeratorOffset::Latest => high_watermark - 1,
496 _ => unreachable!(),
497 };
498 map.insert(*partition, Some(offset));
499 }
500 Ok(map)
501 }
502 KafkaEnumeratorOffset::Timestamp(time) => {
503 let unresolved =
504 sync_resolved_partitions(&mut self.resolved_start_offsets, partitions);
505 if !unresolved.is_empty() {
506 let resolved = self
507 .fetch_offset_for_time(&unresolved, time, watermarks)
508 .await?;
509 self.resolved_start_offsets.extend(resolved);
510 }
511 Ok(self.resolved_start_offsets.clone())
512 }
513 KafkaEnumeratorOffset::None => partitions
514 .iter()
515 .map(|partition| Ok((*partition, None)))
516 .collect(),
517 }
518 }
519
520 async fn fetch_offset_for_time(
521 &self,
522 partitions: &[i32],
523 time: i64,
524 watermarks: &HashMap<i32, (i64, i64)>,
525 ) -> KafkaResult<HashMap<i32, Option<i64>>> {
526 let mut tpl = TopicPartitionList::new();
527
528 for partition in partitions {
529 tpl.add_partition_offset(self.topic.as_str(), *partition, Offset::Offset(time))?;
530 }
531
532 let offsets = self
533 .client
534 .offsets_for_times(tpl, self.sync_call_timeout)
535 .await?;
536
537 let mut result = HashMap::with_capacity(partitions.len());
538
539 for elem in offsets.elements_for_topic(self.topic.as_str()) {
540 match elem.offset() {
541 Offset::Offset(offset) => {
542 result.insert(elem.partition(), Some(offset - 1));
544 }
545 Offset::End => {
546 let (_, high_watermark) = watermarks.get(&elem.partition()).unwrap();
547 tracing::info!(
548 source_id = %self.context.info.source_id,
549 "no message found before timestamp {} (ms) for partition {}, start from latest",
550 time,
551 elem.partition()
552 );
553 result.insert(elem.partition(), Some(high_watermark - 1)); }
555 Offset::Invalid => {
556 tracing::info!(
561 source_id = %self.context.info.source_id,
562 "got invalid offset for partition {} at timestamp {}, align to latest",
563 elem.partition(),
564 time
565 );
566 let (_, high_watermark) = watermarks.get(&elem.partition()).unwrap();
567 result.insert(elem.partition(), Some(high_watermark - 1)); }
569 Offset::Beginning => {
570 let (low, _) = watermarks.get(&elem.partition()).unwrap();
571 tracing::info!(
572 source_id = %self.context.info.source_id,
573 "all message in partition {} is after timestamp {} (ms), start from earliest",
574 elem.partition(),
575 time,
576 );
577 result.insert(elem.partition(), Some(low - 1)); }
579 err_offset @ Offset::Stored | err_offset @ Offset::OffsetTail(_) => {
580 tracing::error!(
581 source_id = %self.context.info.source_id,
582 "got invalid offset for partition {}: {err_offset:?}",
583 elem.partition(),
584 err_offset = err_offset,
585 );
586 return Err(KafkaError::OffsetFetch(RDKafkaErrorCode::NoOffset));
587 }
588 }
589 }
590
591 Ok(result)
592 }
593
594 #[inline]
595 fn report_high_watermark(&mut self, partition: i32, offset: i64) {
596 let high_watermark_metrics =
597 self.high_watermark_metrics
598 .entry(partition)
599 .or_insert_with(|| {
600 self.context
601 .metrics
602 .high_watermark
603 .with_guarded_label_values(&[
604 &self.context.info.source_id.to_string(),
605 &partition.to_string(),
606 ])
607 });
608 high_watermark_metrics.set(offset);
609 }
610
611 pub async fn check_reachability(&self) -> ConnectorResult<()> {
612 let _ = self
613 .client
614 .fetch_metadata(Some(self.topic.as_str()), self.sync_call_timeout)
615 .await?;
616 Ok(())
617 }
618
619 async fn fetch_topic_partition(&self) -> ConnectorResult<Vec<i32>> {
620 let metadata = self
622 .client
623 .fetch_metadata(Some(self.topic.as_str()), self.sync_call_timeout)
624 .await?;
625
626 let topic_meta = match metadata.topics() {
627 [meta] => meta,
628 _ => bail!("topic {} not found", self.topic),
629 };
630
631 if topic_meta.partitions().is_empty() {
632 bail!("topic {} not found", self.topic);
633 }
634
635 Ok(topic_meta
636 .partitions()
637 .iter()
638 .map(|partition| partition.id())
639 .collect())
640 }
641}
642
643fn sync_resolved_partitions(
646 resolved: &mut HashMap<i32, Option<i64>>,
647 partitions: &[i32],
648) -> Vec<i32> {
649 let unchanged = resolved.len() == partitions.len()
652 && partitions
653 .iter()
654 .all(|partition| resolved.contains_key(partition));
655 if unchanged {
656 return Vec::new();
657 }
658 let current: HashSet<i32> = partitions.iter().copied().collect();
659 resolved.retain(|partition, _| current.contains(partition));
660 partitions
661 .iter()
662 .copied()
663 .filter(|partition| !resolved.contains_key(partition))
664 .collect()
665}
666
667#[cfg(test)]
668mod tests {
669 use super::*;
670
671 #[test]
672 fn test_sync_resolved_partitions() {
673 let mut resolved = HashMap::new();
674
675 assert_eq!(sync_resolved_partitions(&mut resolved, &[0, 1]), vec![0, 1]);
677 resolved.extend([(0, Some(10)), (1, Some(20))]);
678
679 assert!(sync_resolved_partitions(&mut resolved, &[0, 1]).is_empty());
681 assert_eq!(resolved, HashMap::from([(0, Some(10)), (1, Some(20))]));
682
683 assert_eq!(sync_resolved_partitions(&mut resolved, &[0, 1, 2]), vec![2]);
685
686 assert!(sync_resolved_partitions(&mut resolved, &[1]).is_empty());
688 assert_eq!(resolved, HashMap::from([(1, Some(20))]));
689 }
690}