1use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17use std::time::{Duration, Instant};
18
19use common_base::readable_size::ReadableSize;
20use common_meta::datanode::TopicStatsReporter;
21use common_meta::distributed_time_constants::TOPIC_STATS_REPORT_INTERVAL_SECS;
22use common_telemetry::{debug, info, warn};
23use common_time::util::current_time_millis;
24use common_wal::config::kafka::DatanodeKafkaConfig;
25use dashmap::DashMap;
26use futures::future::try_join_all;
27use futures_util::StreamExt;
28use rskafka::client::partition::OffsetAt;
29use snafu::{OptionExt, ResultExt};
30use store_api::logstore::entry::{
31 Entry, Id as EntryId, MultiplePartEntry, MultiplePartHeader, NaiveEntry,
32};
33use store_api::logstore::provider::{KafkaProvider, Provider};
34use store_api::logstore::{AppendBatchResponse, LogStore, SendableEntryStream, WalIndex};
35use store_api::storage::RegionId;
36
37use crate::error::{self, ConsumeRecordSnafu, Error, GetOffsetSnafu, InvalidProviderSnafu, Result};
38use crate::kafka::client_manager::{ClientManager, ClientManagerRef};
39use crate::kafka::consumer::{ConsumerBuilder, RecordsBuffer};
40use crate::kafka::index::{
41 GlobalIndexCollector, MIN_BATCH_WINDOW_SIZE, build_region_wal_index_iterator,
42};
43use crate::kafka::periodic_offset_fetcher::PeriodicOffsetFetcher;
44use crate::kafka::producer::OrderedBatchProducerRef;
45use crate::kafka::util::record::{
46 ESTIMATED_META_SIZE, Record, convert_to_kafka_records, maybe_emit_entry, remaining_entries,
47};
48use crate::metrics;
49
50#[derive(Debug, Clone, Copy, Default)]
52pub struct TopicStat {
53 pub latest_offset: u64,
60 pub record_size: u64,
62 pub record_num: u64,
64}
65
66#[derive(Debug)]
68pub struct KafkaLogStore {
69 client_manager: ClientManagerRef,
71 max_batch_bytes: usize,
73 consumer_wait_timeout: Duration,
75 overwrite_entry_start_id: bool,
77 topic_stats: Arc<DashMap<Arc<KafkaProvider>, TopicStat>>,
82}
83
84struct PeriodicTopicStatsReporter {
85 topic_stats: Arc<DashMap<Arc<KafkaProvider>, TopicStat>>,
86 last_reported_timestamp_millis: i64,
87 report_interval_millis: i64,
88}
89
90impl PeriodicTopicStatsReporter {
91 fn align_ts(ts: i64, report_interval_millis: i64) -> i64 {
92 (ts / report_interval_millis) * report_interval_millis
93 }
94
95 fn new(
101 topic_stats: Arc<DashMap<Arc<KafkaProvider>, TopicStat>>,
102 report_interval: Duration,
103 ) -> Self {
104 assert!(!report_interval.is_zero());
105 let report_interval_millis = report_interval.as_millis() as i64;
106 let last_reported_timestamp_millis =
107 Self::align_ts(current_time_millis(), report_interval_millis);
108
109 Self {
110 topic_stats,
111 last_reported_timestamp_millis,
112 report_interval_millis,
113 }
114 }
115}
116
117impl TopicStatsReporter for PeriodicTopicStatsReporter {
118 fn reportable_topics(&mut self) -> Vec<common_meta::datanode::TopicStat> {
119 let now = Self::align_ts(current_time_millis(), self.report_interval_millis);
120 if now < self.last_reported_timestamp_millis + self.report_interval_millis {
121 debug!("Skip reporting topic stats because the interval is not reached");
122 return vec![];
123 }
124
125 self.last_reported_timestamp_millis = now;
126 let mut reportable_topics = Vec::with_capacity(self.topic_stats.len());
127 for e in self.topic_stats.iter() {
128 let topic_stat = e.value();
129 let topic_stat = common_meta::datanode::TopicStat {
130 topic: e.key().topic.clone(),
131 latest_entry_id: topic_stat.latest_offset,
132 record_size: topic_stat.record_size,
133 record_num: topic_stat.record_num,
134 };
135 debug!("Reportable topic: {:?}", topic_stat);
136 reportable_topics.push(topic_stat);
137 }
138 debug!("Reportable {} topics at {}", reportable_topics.len(), now);
139 reportable_topics
140 }
141}
142
143impl KafkaLogStore {
144 pub async fn try_new(
146 config: &DatanodeKafkaConfig,
147 global_index_collector: Option<GlobalIndexCollector>,
148 ) -> Result<Self> {
149 let topic_stats = Arc::new(DashMap::new());
150 let client_manager = Arc::new(
151 ClientManager::try_new(config, global_index_collector, topic_stats.clone()).await?,
152 );
153 let fetcher = PeriodicOffsetFetcher::new(
154 config.topic_latest_offset_fetch_interval,
155 client_manager.clone(),
156 );
157 fetcher.run().await;
158
159 Ok(Self {
160 client_manager,
161 max_batch_bytes: config.max_batch_bytes.as_bytes() as usize,
162 consumer_wait_timeout: config.consumer_wait_timeout,
163 overwrite_entry_start_id: config.overwrite_entry_start_id,
164 topic_stats,
165 })
166 }
167
168 pub fn topic_stats_reporter(&self) -> Box<dyn TopicStatsReporter> {
170 Box::new(PeriodicTopicStatsReporter::new(
171 self.topic_stats.clone(),
172 Duration::from_secs(TOPIC_STATS_REPORT_INTERVAL_SECS),
173 ))
174 }
175}
176
177fn build_entry(
178 data: Vec<u8>,
179 entry_id: EntryId,
180 region_id: RegionId,
181 provider: &Provider,
182 max_data_size: usize,
183) -> Entry {
184 if data.len() <= max_data_size {
185 Entry::Naive(NaiveEntry {
186 provider: provider.clone(),
187 region_id,
188 entry_id,
189 data,
190 })
191 } else {
192 let parts = data
193 .chunks(max_data_size)
194 .map(|s| s.into())
195 .collect::<Vec<_>>();
196 let num_parts = parts.len();
197
198 let mut headers = Vec::with_capacity(num_parts);
199 headers.push(MultiplePartHeader::First);
200 headers.extend((1..num_parts - 1).map(MultiplePartHeader::Middle));
201 headers.push(MultiplePartHeader::Last);
202
203 Entry::MultiplePart(MultiplePartEntry {
204 provider: provider.clone(),
205 region_id,
206 entry_id,
207 headers,
208 parts,
209 })
210 }
211}
212
213#[async_trait::async_trait]
214impl LogStore for KafkaLogStore {
215 type Error = Error;
216
217 fn entry(
219 &self,
220 data: Vec<u8>,
221 entry_id: EntryId,
222 region_id: RegionId,
223 provider: &Provider,
224 ) -> Result<Entry> {
225 provider
226 .as_kafka_provider()
227 .with_context(|| InvalidProviderSnafu {
228 expected: KafkaProvider::type_name(),
229 actual: provider.type_name(),
230 })?;
231
232 let max_data_size = self.max_batch_bytes - ESTIMATED_META_SIZE;
233 Ok(build_entry(
234 data,
235 entry_id,
236 region_id,
237 provider,
238 max_data_size,
239 ))
240 }
241
242 async fn append_batch(&self, entries: Vec<Entry>) -> Result<AppendBatchResponse> {
245 metrics::METRIC_KAFKA_APPEND_BATCH_BYTES_TOTAL.inc_by(
246 entries
247 .iter()
248 .map(|entry| entry.estimated_size())
249 .sum::<usize>() as u64,
250 );
251 let _timer = metrics::METRIC_KAFKA_APPEND_BATCH_ELAPSED.start_timer();
252
253 if entries.is_empty() {
254 return Ok(AppendBatchResponse::default());
255 }
256
257 let region_ids = entries
258 .iter()
259 .map(|entry| entry.region_id())
260 .collect::<HashSet<_>>();
261 let mut region_grouped_records: HashMap<RegionId, (OrderedBatchProducerRef, Vec<_>)> =
262 HashMap::with_capacity(region_ids.len());
263 let mut region_to_provider = HashMap::with_capacity(region_ids.len());
264 for entry in entries {
265 let provider = entry.provider().as_kafka_provider().with_context(|| {
266 error::InvalidProviderSnafu {
267 expected: KafkaProvider::type_name(),
268 actual: entry.provider().type_name(),
269 }
270 })?;
271 region_to_provider.insert(entry.region_id(), provider.clone());
272 let region_id = entry.region_id();
273 match region_grouped_records.entry(region_id) {
274 std::collections::hash_map::Entry::Occupied(mut slot) => {
275 slot.get_mut().1.extend(convert_to_kafka_records(entry)?);
276 }
277 std::collections::hash_map::Entry::Vacant(slot) => {
278 let producer = self
279 .client_manager
280 .get_or_insert(provider)
281 .await?
282 .producer()
283 .clone();
284
285 slot.insert((producer, convert_to_kafka_records(entry)?));
286 }
287 }
288 }
289
290 let mut region_grouped_result_receivers = Vec::with_capacity(region_ids.len());
291 for (region_id, (producer, records)) in region_grouped_records {
292 region_grouped_result_receivers
295 .push((region_id, producer.produce(region_id, records).await?))
296 }
297
298 let region_grouped_max_offset =
299 try_join_all(region_grouped_result_receivers.into_iter().map(
300 |(region_id, receiver)| async move {
301 receiver.wait().await.map(|offset| (region_id, offset))
302 },
303 ))
304 .await?;
305 debug!(
306 "Appended batch to Kafka, region_grouped_max_offset: {:?}",
307 region_grouped_max_offset
308 );
309
310 Ok(AppendBatchResponse {
311 last_entry_ids: region_grouped_max_offset.into_iter().collect(),
312 })
313 }
314
315 async fn read(
318 &self,
319 provider: &Provider,
320 mut entry_id: EntryId,
321 index: Option<WalIndex>,
322 ) -> Result<SendableEntryStream<'static, Entry, Self::Error>> {
323 let provider = provider
324 .as_kafka_provider()
325 .with_context(|| InvalidProviderSnafu {
326 expected: KafkaProvider::type_name(),
327 actual: provider.type_name(),
328 })?;
329
330 let _timer = metrics::METRIC_KAFKA_READ_ELAPSED.start_timer();
331
332 let client = self
334 .client_manager
335 .get_or_insert(provider)
336 .await?
337 .client()
338 .clone();
339
340 if self.overwrite_entry_start_id {
341 let start_offset =
342 client
343 .get_offset(OffsetAt::Earliest)
344 .await
345 .context(GetOffsetSnafu {
346 topic: &provider.topic,
347 })?;
348
349 if entry_id as i64 <= start_offset {
350 warn!(
351 "The entry_id: {} is less than start_offset: {}, topic: {}. Overwriting entry_id with start_offset",
352 entry_id, start_offset, &provider.topic
353 );
354
355 entry_id = start_offset as u64;
356 }
357 }
358
359 let end_offset = client
364 .get_offset(OffsetAt::Latest)
365 .await
366 .context(GetOffsetSnafu {
367 topic: &provider.topic,
368 })?;
369 let latest_offset = (end_offset as u64).saturating_sub(1);
370 self.topic_stats
371 .entry(provider.clone())
372 .and_modify(|stat| {
373 stat.latest_offset = stat.latest_offset.max(latest_offset);
374 })
375 .or_insert_with(|| TopicStat {
376 latest_offset,
377 record_size: 0,
378 record_num: 0,
379 });
380
381 let region_indexes = if let (Some(index), Some(collector)) =
382 (index, self.client_manager.global_index_collector())
383 {
384 collector
385 .read_remote_region_index(index.location_id, provider, index.region_id, entry_id)
386 .await?
387 } else {
388 None
389 };
390
391 let Some(iterator) = build_region_wal_index_iterator(
392 entry_id,
393 end_offset as u64,
394 region_indexes,
395 self.max_batch_bytes,
396 MIN_BATCH_WINDOW_SIZE,
397 ) else {
398 let range = entry_id..end_offset as u64;
399 warn!("No new entries in range {:?} of ns {}", range, provider);
400 return Ok(futures_util::stream::empty().boxed());
401 };
402
403 debug!("Reading entries with {:?} of ns {}", iterator, provider);
404
405 let mut stream_consumer = ConsumerBuilder::default()
407 .client(client)
408 .buffer(RecordsBuffer::new(iterator))
410 .max_batch_size(self.max_batch_bytes)
411 .max_wait_ms(self.consumer_wait_timeout.as_millis() as u32)
412 .build()
413 .unwrap();
414
415 let mut entry_records: HashMap<RegionId, Vec<Record>> = HashMap::new();
417 let provider = provider.clone();
418 let stream = async_stream::stream!({
419 let now = Instant::now();
420 while let Some(consume_result) = stream_consumer.next().await {
421 let (record_and_offset, high_watermark) =
425 consume_result.context(ConsumeRecordSnafu {
426 topic: &provider.topic,
427 })?;
428 let (kafka_record, offset) = (record_and_offset.record, record_and_offset.offset);
429
430 debug!(
431 "Read a record at offset {} for topic {}, high watermark: {}",
432 offset, provider.topic, high_watermark
433 );
434
435 if kafka_record.value.is_none() {
437 if check_termination(offset, end_offset) {
438 if let Some(entries) = remaining_entries(&provider, &mut entry_records) {
439 yield Ok(entries);
440 }
441 break;
442 }
443 continue;
444 }
445
446 let record = Record::try_from(kafka_record)?;
447 if let Some(mut entry) = maybe_emit_entry(&provider, record, &mut entry_records)? {
449 entry.set_entry_id(offset as u64);
453 yield Ok(vec![entry]);
454 }
455
456 if check_termination(offset, end_offset) {
457 if let Some(entries) = remaining_entries(&provider, &mut entry_records) {
458 yield Ok(entries);
459 }
460 break;
461 }
462 }
463
464 metrics::METRIC_KAFKA_READ_BYTES_TOTAL.inc_by(stream_consumer.total_fetched_bytes());
465
466 info!(
467 "Fetched {} bytes from topic: {}, start_entry_id: {}, end_offset: {}, elapsed: {:?}",
468 ReadableSize(stream_consumer.total_fetched_bytes()),
469 stream_consumer.topic(),
470 entry_id,
471 end_offset,
472 now.elapsed()
473 );
474 });
475 Ok(Box::pin(stream))
476 }
477
478 async fn create_namespace(&self, _provider: &Provider) -> Result<()> {
480 Ok(())
481 }
482
483 async fn delete_namespace(&self, _provider: &Provider) -> Result<()> {
485 Ok(())
486 }
487
488 async fn list_namespaces(&self) -> Result<Vec<Provider>> {
490 Ok(vec![])
491 }
492
493 async fn obsolete(
497 &self,
498 provider: &Provider,
499 region_id: RegionId,
500 entry_id: EntryId,
501 ) -> Result<()> {
502 if let Some(collector) = self.client_manager.global_index_collector() {
503 let provider = provider
504 .as_kafka_provider()
505 .with_context(|| InvalidProviderSnafu {
506 expected: KafkaProvider::type_name(),
507 actual: provider.type_name(),
508 })?;
509 collector.truncate(provider, region_id, entry_id).await?;
510 }
511 Ok(())
512 }
513
514 async fn obsolete_all(&self, _provider: &Provider, _region_id: RegionId) -> Result<()> {
515 Ok(())
516 }
517
518 fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId> {
520 let provider = provider
521 .as_kafka_provider()
522 .with_context(|| InvalidProviderSnafu {
523 expected: KafkaProvider::type_name(),
524 actual: provider.type_name(),
525 })?;
526
527 let stat = self
528 .topic_stats
529 .get(provider)
530 .as_deref()
531 .copied()
532 .unwrap_or_default();
533
534 Ok(stat.latest_offset)
535 }
536
537 async fn stop(&self) -> Result<()> {
539 Ok(())
540 }
541}
542
543fn check_termination(offset: i64, end_offset: i64) -> bool {
544 if offset >= end_offset {
546 debug!("Stream consumer terminates at offset {}", offset);
547 true
549 } else {
550 false
551 }
552}
553
554#[cfg(test)]
555mod tests {
556
557 use std::assert_matches;
558 use std::collections::HashMap;
559 use std::sync::Arc;
560 use std::time::Duration;
561
562 use common_base::readable_size::ReadableSize;
563 use common_meta::datanode::TopicStatsReporter;
564 use common_telemetry::info;
565 use common_telemetry::tracing::warn;
566 use common_wal::config::kafka::DatanodeKafkaConfig;
567 use common_wal::config::kafka::common::KafkaConnectionConfig;
568 use dashmap::DashMap;
569 use futures::TryStreamExt;
570 use rand::Rng;
571 use rand::prelude::SliceRandom;
572 use rskafka::client::partition::OffsetAt;
573 use store_api::logstore::LogStore;
574 use store_api::logstore::entry::{Entry, MultiplePartEntry, MultiplePartHeader, NaiveEntry};
575 use store_api::logstore::provider::Provider;
576 use store_api::storage::RegionId;
577
578 use super::build_entry;
579 use crate::kafka::log_store::{KafkaLogStore, PeriodicTopicStatsReporter, TopicStat};
580
581 #[test]
582 fn test_build_naive_entry() {
583 let provider = Provider::kafka_provider("my_topic".to_string());
584 let region_id = RegionId::new(1, 1);
585 let entry = build_entry(vec![1; 100], 1, region_id, &provider, 120);
586
587 assert_eq!(
588 entry.into_naive_entry().unwrap(),
589 NaiveEntry {
590 provider,
591 region_id,
592 entry_id: 1,
593 data: vec![1; 100]
594 }
595 )
596 }
597
598 #[test]
599 fn test_build_into_multiple_part_entry() {
600 let provider = Provider::kafka_provider("my_topic".to_string());
601 let region_id = RegionId::new(1, 1);
602 let entry = build_entry(vec![1; 100], 1, region_id, &provider, 50);
603
604 assert_eq!(
605 entry.into_multiple_part_entry().unwrap(),
606 MultiplePartEntry {
607 provider: provider.clone(),
608 region_id,
609 entry_id: 1,
610 headers: vec![MultiplePartHeader::First, MultiplePartHeader::Last],
611 parts: vec![vec![1; 50], vec![1; 50]],
612 }
613 );
614
615 let region_id = RegionId::new(1, 1);
616 let entry = build_entry(vec![1; 100], 1, region_id, &provider, 21);
617
618 assert_eq!(
619 entry.into_multiple_part_entry().unwrap(),
620 MultiplePartEntry {
621 provider,
622 region_id,
623 entry_id: 1,
624 headers: vec![
625 MultiplePartHeader::First,
626 MultiplePartHeader::Middle(1),
627 MultiplePartHeader::Middle(2),
628 MultiplePartHeader::Middle(3),
629 MultiplePartHeader::Last
630 ],
631 parts: vec![
632 vec![1; 21],
633 vec![1; 21],
634 vec![1; 21],
635 vec![1; 21],
636 vec![1; 16]
637 ],
638 }
639 )
640 }
641
642 fn generate_entries(
643 logstore: &KafkaLogStore,
644 provider: &Provider,
645 num_entries: usize,
646 region_id: RegionId,
647 data_len: usize,
648 ) -> Vec<Entry> {
649 (0..num_entries)
650 .map(|_| {
651 let data: Vec<u8> = (0..data_len).map(|_| rand::random::<u8>()).collect();
652 logstore.entry(data, 0, region_id, provider).unwrap()
654 })
655 .collect()
656 }
657
658 async fn prepare_topic(logstore: &KafkaLogStore, topic_name: &str) {
659 let controller_client = logstore.client_manager.controller_client();
660 controller_client
661 .create_topic(topic_name.to_string(), 1, 1, 5000)
662 .await
663 .unwrap();
664 }
665
666 #[tokio::test]
667 async fn test_append_batch_basic() {
668 common_telemetry::init_default_ut_logging();
669 let Ok(broker_endpoints) = std::env::var("GT_KAFKA_ENDPOINTS") else {
670 warn!("The endpoints is empty, skipping the test 'test_append_batch_basic'");
671 return;
672 };
673 let broker_endpoints = broker_endpoints
674 .split(',')
675 .map(|s| s.trim().to_string())
676 .collect::<Vec<_>>();
677 let config = DatanodeKafkaConfig {
678 connection: KafkaConnectionConfig {
679 broker_endpoints,
680 ..Default::default()
681 },
682 max_batch_bytes: ReadableSize::kb(32),
683 ..Default::default()
684 };
685 let logstore = KafkaLogStore::try_new(&config, None).await.unwrap();
686 let topic_name = uuid::Uuid::new_v4().to_string();
687 prepare_topic(&logstore, &topic_name).await;
688 let provider = Provider::kafka_provider(topic_name);
689
690 let region_entries = (0..5)
691 .map(|i| {
692 let region_id = RegionId::new(1, i);
693 (
694 region_id,
695 generate_entries(&logstore, &provider, 20, region_id, 1024),
696 )
697 })
698 .collect::<HashMap<RegionId, Vec<_>>>();
699
700 let mut all_entries = region_entries
701 .values()
702 .flatten()
703 .cloned()
704 .collect::<Vec<_>>();
705 all_entries.shuffle(&mut rand::rng());
706
707 let response = logstore.append_batch(all_entries.clone()).await.unwrap();
708 assert_eq!(response.last_entry_ids.len(), 5);
710 let got_entries = logstore
711 .read(&provider, 0, None)
712 .await
713 .unwrap()
714 .try_collect::<Vec<_>>()
715 .await
716 .unwrap()
717 .into_iter()
718 .flatten()
719 .collect::<Vec<_>>();
720 for (region_id, _) in region_entries {
721 let expected_entries = all_entries
722 .iter()
723 .filter(|entry| entry.region_id() == region_id)
724 .cloned()
725 .collect::<Vec<_>>();
726 let mut actual_entries = got_entries
727 .iter()
728 .filter(|entry| entry.region_id() == region_id)
729 .cloned()
730 .collect::<Vec<_>>();
731 actual_entries
732 .iter_mut()
733 .for_each(|entry| entry.set_entry_id(0));
734 assert_eq!(expected_entries, actual_entries);
735 }
736 let latest_entry_id = logstore.latest_entry_id(&provider).unwrap();
737 let client = logstore
738 .client_manager
739 .get_or_insert(provider.as_kafka_provider().unwrap())
740 .await
741 .unwrap();
742 assert_eq!(latest_entry_id, 99);
743 let latest = client.client().get_offset(OffsetAt::Latest).await.unwrap();
745 assert_eq!(latest, 100);
746 }
747
748 #[tokio::test]
749 async fn test_append_batch_basic_large() {
750 common_telemetry::init_default_ut_logging();
751 let Ok(broker_endpoints) = std::env::var("GT_KAFKA_ENDPOINTS") else {
752 warn!("The endpoints is empty, skipping the test 'test_append_batch_basic_large'");
753 return;
754 };
755 let data_size_kb = rand::rng().random_range(9..31usize);
756 info!("Entry size: {}Ki", data_size_kb);
757 let broker_endpoints = broker_endpoints
758 .split(',')
759 .map(|s| s.trim().to_string())
760 .collect::<Vec<_>>();
761 let config = DatanodeKafkaConfig {
762 connection: KafkaConnectionConfig {
763 broker_endpoints,
764 ..Default::default()
765 },
766 max_batch_bytes: ReadableSize::kb(8),
767 ..Default::default()
768 };
769 let logstore = KafkaLogStore::try_new(&config, None).await.unwrap();
770 let topic_name = uuid::Uuid::new_v4().to_string();
771 prepare_topic(&logstore, &topic_name).await;
772 let provider = Provider::kafka_provider(topic_name);
773 let region_entries = (0..5)
774 .map(|i| {
775 let region_id = RegionId::new(1, i);
776 (
777 region_id,
778 generate_entries(&logstore, &provider, 20, region_id, data_size_kb * 1024),
779 )
780 })
781 .collect::<HashMap<RegionId, Vec<_>>>();
782
783 let mut all_entries = region_entries
784 .values()
785 .flatten()
786 .cloned()
787 .collect::<Vec<_>>();
788 assert_matches!(all_entries[0], Entry::MultiplePart(_));
789 all_entries.shuffle(&mut rand::rng());
790
791 let response = logstore.append_batch(all_entries.clone()).await.unwrap();
792 assert_eq!(response.last_entry_ids.len(), 5);
794 let got_entries = logstore
795 .read(&provider, 0, None)
796 .await
797 .unwrap()
798 .try_collect::<Vec<_>>()
799 .await
800 .unwrap()
801 .into_iter()
802 .flatten()
803 .collect::<Vec<_>>();
804 for (region_id, _) in region_entries {
805 let expected_entries = all_entries
806 .iter()
807 .filter(|entry| entry.region_id() == region_id)
808 .cloned()
809 .collect::<Vec<_>>();
810 let mut actual_entries = got_entries
811 .iter()
812 .filter(|entry| entry.region_id() == region_id)
813 .cloned()
814 .collect::<Vec<_>>();
815 actual_entries
816 .iter_mut()
817 .for_each(|entry| entry.set_entry_id(0));
818 assert_eq!(expected_entries, actual_entries);
819 }
820 let high_wathermark = logstore.latest_entry_id(&provider).unwrap();
821 assert_eq!(high_wathermark, (data_size_kb as u64 / 8 + 1) * 20 * 5 - 1);
822 }
823
824 #[tokio::test]
825 async fn test_topic_stats_reporter() {
826 common_telemetry::init_default_ut_logging();
827 let topic_stats = Arc::new(DashMap::new());
828 let provider = Provider::kafka_provider("my_topic".to_string());
829 topic_stats.insert(
830 provider.as_kafka_provider().unwrap().clone(),
831 TopicStat {
832 latest_offset: 0,
833 record_size: 0,
834 record_num: 0,
835 },
836 );
837 let mut reporter = PeriodicTopicStatsReporter::new(topic_stats, Duration::from_secs(1));
838 let reportable_topics = reporter.reportable_topics();
840 assert_eq!(reportable_topics.len(), 0);
841
842 tokio::time::sleep(Duration::from_secs(1)).await;
844 let reportable_topics = reporter.reportable_topics();
845 assert_eq!(reportable_topics.len(), 1);
846
847 let reportable_topics = reporter.reportable_topics();
849 assert_eq!(reportable_topics.len(), 0);
850 }
851}