1use std::collections::{BTreeMap, BTreeSet, HashMap};
16use std::mem;
17use std::sync::Arc;
18use std::time::SystemTime;
19
20use await_tree::{InstrumentAwait, SpanExt};
21use bytes::{Bytes, BytesMut};
22use risingwave_common::catalog::TableId;
23use risingwave_common::hash::VirtualNode;
24use risingwave_common::util::row_serde::OrderedRowSerde;
25use risingwave_hummock_sdk::key::{FullKey, MAX_KEY_LEN, TABLE_PREFIX_LEN, UserKey, user_key};
26use risingwave_hummock_sdk::key_range::KeyRange;
27use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner, VnodeStatistics};
28use risingwave_hummock_sdk::table_stats::{TableStats, TableStatsMap};
29use risingwave_hummock_sdk::{HummockEpoch, HummockSstableObjectId, LocalSstableInfo};
30use risingwave_pb::hummock::{PbSstableFilterLayout, PbSstableFilterType};
31
32use super::utils::CompressionAlgorithm;
33use super::{
34 BlockBuilder, BlockBuilderOptions, BlockMeta, DEFAULT_BLOCK_SIZE, DEFAULT_ENTRY_SIZE,
35 DEFAULT_RESTART_INTERVAL, SstableMeta, SstableWriter, VERSION,
36};
37use crate::compaction_catalog_manager::{
38 CompactionCatalogAgent, CompactionCatalogAgentRef, FilterKeyExtractorImpl,
39 FullKeyFilterKeyExtractor,
40};
41use crate::hummock::sstable::{
42 DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP, FilterBuilder, FilterBuilderOptions, utils,
43};
44use crate::hummock::value::HummockValue;
45use crate::hummock::{
46 Block, BlockHolder, BlockIterator, HummockResult, MemoryLimiter, Xor16FilterBuilder,
47 try_shorten_block_smallest_key,
48};
49use crate::monitor::CompactorMetrics;
50use crate::opts::StorageOpts;
51
52pub const DEFAULT_SSTABLE_SIZE: usize = 4 * 1024 * 1024;
53pub const DEFAULT_BLOOM_FALSE_POSITIVE: f64 = 0.001;
54pub const DEFAULT_MAX_SST_SIZE: u64 = 512 * 1024 * 1024;
55pub const MIN_BLOCK_SIZE: usize = 8 * 1024;
56
57#[derive(Clone, Debug)]
58pub struct SstableBuilderOptions {
59 pub capacity: usize,
61 pub block_capacity: usize,
63 pub restart_interval: usize,
65 pub bloom_false_positive: f64,
67 pub compression_algorithm: CompressionAlgorithm,
69 pub max_sst_size: u64,
70 pub shorten_block_meta_key_threshold: Option<usize>,
72 pub max_vnode_key_range_bytes: Option<usize>,
74 pub estimated_output_key_count: Option<usize>,
76 pub filter_hash_prealloc_key_count_cap: usize,
78}
79
80impl From<&StorageOpts> for SstableBuilderOptions {
81 fn from(options: &StorageOpts) -> SstableBuilderOptions {
82 let capacity: usize = (options.sstable_size_mb as usize) * (1 << 20);
83 SstableBuilderOptions {
84 capacity,
85 block_capacity: (options.block_size_kb as usize) * (1 << 10),
86 restart_interval: DEFAULT_RESTART_INTERVAL,
87 bloom_false_positive: options.bloom_false_positive,
88 compression_algorithm: CompressionAlgorithm::None,
89 max_sst_size: options.compactor_max_sst_size,
90 shorten_block_meta_key_threshold: options.shorten_block_meta_key_threshold,
91 max_vnode_key_range_bytes: None,
92 estimated_output_key_count: None,
93 filter_hash_prealloc_key_count_cap: DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP,
94 }
95 }
96}
97
98impl Default for SstableBuilderOptions {
99 fn default() -> Self {
100 Self {
101 capacity: DEFAULT_SSTABLE_SIZE,
102 block_capacity: DEFAULT_BLOCK_SIZE,
103 restart_interval: DEFAULT_RESTART_INTERVAL,
104 bloom_false_positive: DEFAULT_BLOOM_FALSE_POSITIVE,
105 compression_algorithm: CompressionAlgorithm::None,
106 max_sst_size: DEFAULT_MAX_SST_SIZE,
107 shorten_block_meta_key_threshold: None,
108 max_vnode_key_range_bytes: None,
109 estimated_output_key_count: None,
110 filter_hash_prealloc_key_count_cap: DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP,
111 }
112 }
113}
114
115impl SstableBuilderOptions {
116 pub fn estimated_output_key_count(&self) -> usize {
117 self.estimated_output_key_count
118 .unwrap_or(self.capacity / DEFAULT_ENTRY_SIZE + 1)
119 }
120
121 pub fn filter_builder_options(&self) -> FilterBuilderOptions {
122 FilterBuilderOptions {
123 estimated_key_count: self.estimated_output_key_count(),
124 estimated_block_count: self.capacity / self.block_capacity + 1,
125 hash_prealloc_key_count_cap: self.filter_hash_prealloc_key_count_cap,
126 }
127 }
128}
129
130pub struct SstableBuilderOutput<WO> {
131 pub sst_info: LocalSstableInfo,
132 pub writer_output: WO,
133 pub stats: SstableBuilderOutputStats,
134}
135
136pub struct SstableBuilder<W: SstableWriter, F: FilterBuilder> {
137 options: SstableBuilderOptions,
139 writer: W,
141 block_builder: BlockBuilder,
143
144 compaction_catalog_agent_ref: CompactionCatalogAgentRef,
145 block_metas: Vec<BlockMeta>,
147
148 table_ids: BTreeSet<TableId>,
150 last_full_key: Vec<u8>,
151 raw_key: Vec<u8>,
153 raw_value: BytesMut,
155 last_table_id: Option<TableId>,
156 sst_object_id: HummockSstableObjectId,
157
158 table_stats: TableStatsMap,
160 last_table_stats: TableStats,
163
164 filter_builder: F,
165
166 epoch_set: BTreeSet<u64>,
167 memory_limiter: Option<Arc<MemoryLimiter>>,
168
169 block_size_vec: Vec<usize>, vnode_range_collector: Option<VnodeUserKeyRangeCollector>,
171}
172
173impl<W: SstableWriter> SstableBuilder<W, Xor16FilterBuilder> {
174 pub fn for_test(
175 sstable_id: u64,
176 writer: W,
177 options: SstableBuilderOptions,
178 table_id_to_vnode: HashMap<impl Into<TableId>, usize>,
179 table_id_to_watermark_serde: HashMap<
180 impl Into<TableId>,
181 Option<(OrderedRowSerde, OrderedRowSerde, usize)>,
182 >,
183 ) -> Self {
184 let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
185 FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
186 table_id_to_vnode
187 .into_iter()
188 .map(|(table_id, v)| (table_id.into(), v))
189 .collect(),
190 table_id_to_watermark_serde
191 .into_iter()
192 .map(|(table_id, v)| (table_id.into(), v))
193 .collect(),
194 HashMap::default(),
195 ));
196
197 Self::new(
198 sstable_id,
199 writer,
200 Xor16FilterBuilder::create(options.filter_builder_options()),
201 options,
202 compaction_catalog_agent_ref,
203 None,
204 )
205 }
206}
207
208impl<W: SstableWriter, F: FilterBuilder> SstableBuilder<W, F> {
209 pub fn new(
210 sst_object_id: impl Into<HummockSstableObjectId>,
211 writer: W,
212 filter_builder: F,
213 options: SstableBuilderOptions,
214 compaction_catalog_agent_ref: CompactionCatalogAgentRef,
215 memory_limiter: Option<Arc<MemoryLimiter>>,
216 ) -> Self {
217 let sst_object_id = sst_object_id.into();
218 Self {
219 vnode_range_collector: VnodeUserKeyRangeCollector::with_limit(
220 options.max_vnode_key_range_bytes,
221 ),
222 options: options.clone(),
223 writer,
224 block_builder: BlockBuilder::new(BlockBuilderOptions {
225 capacity: options.block_capacity,
226 restart_interval: options.restart_interval,
227 compression_algorithm: options.compression_algorithm,
228 }),
229 filter_builder,
230 block_metas: Vec::with_capacity(options.capacity / options.block_capacity + 1),
231 table_ids: BTreeSet::new(),
232 last_table_id: None,
233 raw_key: Vec::new(),
234 raw_value: BytesMut::new(),
235 last_full_key: vec![],
236 sst_object_id,
237 compaction_catalog_agent_ref,
238 table_stats: Default::default(),
239 last_table_stats: Default::default(),
240 epoch_set: BTreeSet::default(),
241 memory_limiter,
242 block_size_vec: Vec::new(),
243 }
244 }
245
246 pub async fn add_for_test(
248 &mut self,
249 full_key: FullKey<&[u8]>,
250 value: HummockValue<&[u8]>,
251 ) -> HummockResult<()> {
252 self.add(full_key, value).await
253 }
254
255 pub fn current_block_size(&self) -> usize {
256 self.block_builder.approximate_len()
257 }
258
259 pub async fn add_raw_block(
262 &mut self,
263 mut buf: Bytes,
264 filter_data: Vec<u8>,
265 smallest_key: FullKey<Vec<u8>>,
266 largest_key: Vec<u8>,
267 mut meta: BlockMeta,
268 ) -> HummockResult<Option<usize>> {
269 let table_id = smallest_key.user_key.table_id;
270 if self.last_table_id.is_none() || self.last_table_id.unwrap() != table_id {
271 if !self.block_builder.is_empty() {
272 self.build_block()
274 .instrument_await("sstable_build_block_on_raw_table_switch".verbose())
275 .await?;
276 }
277
278 self.table_ids.insert(table_id);
279 self.finalize_last_table_stats();
280 self.last_table_id = Some(table_id);
281 }
282
283 if !self.block_builder.is_empty() {
284 let min_block_size = std::cmp::min(MIN_BLOCK_SIZE, self.options.block_capacity / 4);
285
286 if self.block_builder.approximate_len() < min_block_size {
288 let block = Block::decode(buf, meta.uncompressed_size as usize)?;
289 let mut iter = BlockIterator::new(BlockHolder::from_owned_block(Box::new(block)));
290 iter.seek_to_first();
291 while iter.is_valid() {
292 let value = HummockValue::from_slice(iter.value()).unwrap_or_else(|_| {
293 panic!(
294 "decode failed for fast compact sst_id {} block_idx {} last_table_id {:?}",
295 self.sst_object_id, self.block_metas.len(), self.last_table_id
296 )
297 });
298 self.add_impl(iter.key(), value, false)
299 .instrument_await("sstable_add_raw_block_fallback_add".verbose())
300 .await?;
301 iter.next();
302 }
303 return Ok(None);
304 }
305
306 self.build_block()
307 .instrument_await("sstable_build_block_before_raw_block".verbose())
308 .await?;
309 }
310 let algorithm = Block::get_algorithm(&buf)?;
313 if algorithm == CompressionAlgorithm::None
314 && algorithm != self.options.compression_algorithm
315 {
316 buf = BlockBuilder::compress_block(buf, self.options.compression_algorithm)?;
317 meta.len = buf.len() as u32;
318 }
319 self.last_full_key = largest_key;
320 assert_eq!(
321 meta.len as usize,
322 buf.len(),
323 "meta {} buf {} last_table_id {:?}",
324 meta.len,
325 buf.len(),
326 self.last_table_id
327 );
328 meta.offset = self.writer.data_len() as u32;
329 let block_len = buf.len();
330 self.block_metas.push(meta);
331 self.filter_builder.add_raw_data(filter_data);
332 let block_meta = self.block_metas.last_mut().unwrap();
333 self.writer
334 .write_block_bytes(buf, block_meta)
335 .instrument_await("sstable_write_raw_block_bytes".verbose())
336 .await?;
337
338 Ok(Some(block_len))
339 }
340
341 pub async fn add(
343 &mut self,
344 full_key: FullKey<&[u8]>,
345 value: HummockValue<&[u8]>,
346 ) -> HummockResult<()> {
347 self.add_impl(full_key, value, true).await
348 }
349
350 async fn add_impl(
352 &mut self,
353 full_key: FullKey<&[u8]>,
354 value: HummockValue<&[u8]>,
355 could_switch_block: bool,
356 ) -> HummockResult<()> {
357 const LARGE_KEY_LEN: usize = MAX_KEY_LEN >> 1;
358
359 let table_key_len = full_key.user_key.table_key.as_ref().len();
360 let table_value_len = match &value {
361 HummockValue::Put(t) => t.len(),
362 HummockValue::Delete => 0,
363 };
364 let large_value_len = self.options.max_sst_size as usize / 10;
365 let large_key_value_len = self.options.max_sst_size as usize / 2;
366 if table_key_len >= LARGE_KEY_LEN
367 || table_value_len > large_value_len
368 || table_key_len + table_value_len > large_key_value_len
369 {
370 let table_id = full_key.user_key.table_id;
371 tracing::warn!(
372 "A large key/value (table_id={}, key len={}, value len={}, epoch={}, spill offset={}) is added to block",
373 table_id,
374 table_key_len,
375 table_value_len,
376 full_key.epoch_with_gap.pure_epoch(),
377 full_key.epoch_with_gap.offset(),
378 );
379 }
380
381 full_key.encode_into(&mut self.raw_key);
383 value.encode(&mut self.raw_value);
384 let is_new_user_key = self.last_full_key.is_empty()
385 || !user_key(&self.raw_key).eq(user_key(self.last_full_key.as_slice()));
386 let table_id = full_key.user_key.table_id;
387 let is_new_table = self.last_table_id != Some(table_id);
388 let current_block_size = self.current_block_size();
389 let is_block_full = current_block_size >= self.options.block_capacity
390 || (current_block_size > self.options.block_capacity / 4 * 3
391 && current_block_size + self.raw_value.len() + self.raw_key.len()
392 > self.options.block_capacity);
393
394 if is_new_table {
395 assert!(
396 could_switch_block,
397 "is_new_user_key {} sst_id {} block_idx {} table_id {} last_table_id {:?} full_key {:?}",
398 is_new_user_key,
399 self.sst_object_id,
400 self.block_metas.len(),
401 table_id,
402 self.last_table_id,
403 full_key
404 );
405 self.table_ids.insert(table_id);
406 self.finalize_last_table_stats();
407 self.last_table_id = Some(table_id);
408 if !self.block_builder.is_empty() {
409 self.build_block()
410 .instrument_await("sstable_build_block_on_table_switch".verbose())
411 .await?;
412 }
413 } else if is_block_full && could_switch_block {
414 self.build_block()
415 .instrument_await("sstable_build_block_on_block_full".verbose())
416 .await?;
417 }
418 self.last_table_stats.total_key_count += 1;
419 self.epoch_set.insert(full_key.epoch_with_gap.pure_epoch());
420
421 if self.block_builder.is_empty() {
423 let smallest_key = if let Some(threshold) =
424 self.options.shorten_block_meta_key_threshold
425 && !self.last_full_key.is_empty()
426 && full_key.encoded_len() >= threshold
427 {
428 let prev = FullKey::decode(&self.last_full_key);
429 if let Some(shortened) = try_shorten_block_smallest_key(&prev, &full_key) {
430 shortened.encode()
431 } else {
432 full_key.encode()
433 }
434 } else {
435 full_key.encode()
436 };
437
438 self.block_metas.push(BlockMeta {
439 offset: utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
440 panic!(
441 "WARN overflow can't convert writer_data_len {} into u32 sst_id {} block_idx {} tables {:?}",
442 self.writer.data_len(),
443 self.sst_object_id,
444 self.block_metas.len(),
445 self.table_ids,
446 )
447 }),
448 len: 0,
449 smallest_key,
450 uncompressed_size: 0,
451 total_key_count: 0,
452 stale_key_count: 0,
453 });
454 }
455
456 let filter_key = self
457 .compaction_catalog_agent_ref
458 .extract(user_key(&self.raw_key));
459 if !filter_key.is_empty() {
461 self.filter_builder
462 .add_key(filter_key, table_id.as_raw_id());
463 }
464 self.block_builder.add(
466 table_id,
467 &self.raw_key[TABLE_PREFIX_LEN..],
468 self.raw_value.as_ref(),
469 );
470 self.block_metas.last_mut().unwrap().total_key_count += 1;
471 if !is_new_user_key || value.is_delete() {
472 self.block_metas.last_mut().unwrap().stale_key_count += 1;
473 }
474 self.last_table_stats.total_key_size += full_key.encoded_len() as i64;
475 self.last_table_stats.total_value_size += value.encoded_len() as i64;
476
477 if let Some(collector) = self.vnode_range_collector.as_mut() {
478 collector.observe_key(
479 VirtualNode::from_index(full_key.user_key.get_vnode_id()),
480 &self.raw_key,
481 &self.last_full_key,
482 );
483 }
484
485 mem::swap(&mut self.last_full_key, &mut self.raw_key);
486
487 self.raw_key.clear();
488 self.raw_value.clear();
489 Ok(())
490 }
491
492 pub async fn finish(mut self) -> HummockResult<SstableBuilderOutput<W::Output>> {
502 let smallest_key = if self.block_metas.is_empty() {
503 vec![]
504 } else {
505 self.block_metas[0].smallest_key.clone()
506 };
507 let largest_key = self.last_full_key.clone();
508 self.finalize_last_table_stats();
509
510 assert!(
513 self.table_ids.len() <= 1 || self.vnode_range_collector.is_none(),
514 "vnode key-range hints are only supported for single-table SSTs, found {} tables",
515 self.table_ids.len()
516 );
517
518 self.build_block()
519 .instrument_await("sstable_finish_build_block".verbose())
520 .await?;
521 let right_exclusive = false;
522 let meta_offset = self.writer.data_len() as u64;
523
524 let filter_data = self.filter_builder.finish(self.memory_limiter.clone());
525 let (filter_type, filter_layout) = if filter_data.is_none() {
526 (
527 PbSstableFilterType::SstableFilterNone,
528 PbSstableFilterLayout::Unspecified,
529 )
530 } else if self.filter_builder.support_blocked_raw_data() {
531 (
532 self.filter_builder.filter_type(),
533 PbSstableFilterLayout::Blocked,
534 )
535 } else {
536 (
537 self.filter_builder.filter_type(),
538 PbSstableFilterLayout::Plain,
539 )
540 };
541
542 let total_key_count = self
543 .block_metas
544 .iter()
545 .map(|block_meta| block_meta.total_key_count as u64)
546 .sum::<u64>();
547 let stale_key_count = self
548 .block_metas
549 .iter()
550 .map(|block_meta| block_meta.stale_key_count as u64)
551 .sum::<u64>();
552 let uncompressed_file_size = self
553 .block_metas
554 .iter()
555 .map(|block_meta| block_meta.uncompressed_size as u64)
556 .sum::<u64>();
557
558 #[expect(deprecated)]
559 let mut meta = SstableMeta {
560 block_metas: self.block_metas,
561 bloom_filter: filter_data.unwrap_or_default(),
562 estimated_size: 0,
563 key_count: utils::checked_into_u32(total_key_count).unwrap_or_else(|_| {
564 panic!(
565 "WARN overflow can't convert total_key_count {} into u32 tables {:?}",
566 total_key_count, self.table_ids,
567 )
568 }),
569 smallest_key,
570 largest_key,
571 version: VERSION,
572 meta_offset,
573 monotonic_tombstone_events: vec![],
574 };
575
576 let meta_encode_size = meta.encoded_size();
577 let encoded_size_u32 = utils::checked_into_u32(meta_encode_size).unwrap_or_else(|_| {
578 panic!(
579 "WARN overflow can't convert meta_encoded_size {} into u32 tables {:?}",
580 meta_encode_size, self.table_ids,
581 )
582 });
583 let meta_offset_u32 = utils::checked_into_u32(meta_offset).unwrap_or_else(|_| {
584 panic!(
585 "WARN overflow can't convert meta_offset {} into u32 tables {:?}",
586 meta_offset, self.table_ids,
587 )
588 });
589 meta.estimated_size = encoded_size_u32
590 .checked_add(meta_offset_u32)
591 .unwrap_or_else(|| {
592 panic!(
593 "WARN overflow encoded_size_u32 {} meta_offset_u32 {} table_id {:?} table_ids {:?}",
594 encoded_size_u32, meta_offset_u32, self.last_table_id, self.table_ids
595 )
596 });
597
598 let (avg_key_size, avg_value_size) = if self.table_stats.is_empty() {
599 (0, 0)
600 } else {
601 let total_key_count: usize = self
602 .table_stats
603 .values()
604 .map(|s| s.total_key_count as usize)
605 .sum();
606
607 let total_key_size: usize = self
608 .table_stats
609 .values()
610 .map(|s| s.total_key_size as usize)
611 .sum();
612
613 let total_value_size: usize = self
614 .table_stats
615 .values()
616 .map(|s| s.total_value_size as usize)
617 .sum();
618
619 (
620 total_key_size.checked_div(total_key_count).unwrap_or(0),
621 total_value_size.checked_div(total_key_count).unwrap_or(0),
622 )
623 };
624
625 let (min_epoch, max_epoch) = {
626 if self.epoch_set.is_empty() {
627 (HummockEpoch::MAX, u64::MIN)
628 } else {
629 (
630 *self.epoch_set.first().unwrap(),
631 *self.epoch_set.last().unwrap(),
632 )
633 }
634 };
635
636 let vnode_user_key_ranges = self
637 .vnode_range_collector
638 .take()
639 .and_then(|collector| collector.finish(&self.last_full_key));
640
641 let sst_info: SstableInfo = SstableInfoInner {
642 object_id: self.sst_object_id,
643 sst_id: self.sst_object_id.as_raw_id().into(),
645 key_range: KeyRange {
646 left: Bytes::from(meta.smallest_key.clone()),
647 right: Bytes::from(meta.largest_key.clone()),
648 right_exclusive,
649 },
650 file_size: meta.estimated_size as u64,
651 table_ids: self.table_ids.into_iter().collect(),
652 meta_offset: meta.meta_offset,
653 stale_key_count,
654 total_key_count,
655 uncompressed_file_size: uncompressed_file_size + meta.encoded_size() as u64,
656 min_epoch,
657 max_epoch,
658 range_tombstone_count: 0,
659 sst_size: meta.estimated_size as u64,
660 filter_type,
661 filter_layout,
662 vnode_statistics: vnode_user_key_ranges,
663 }
664 .into();
665
666 tracing::trace!(
667 "meta_size {} filter_size {} add_key_counts {} stale_key_count {} min_epoch {} max_epoch {} epoch_count {}",
668 meta.encoded_size(),
669 meta.bloom_filter.len(),
670 total_key_count,
671 stale_key_count,
672 min_epoch,
673 max_epoch,
674 self.epoch_set.len()
675 );
676 let filter_size = meta.bloom_filter.len();
677 let sstable_file_size = sst_info.file_size as usize;
678
679 if !meta.block_metas.is_empty() {
680 let mut last_table_id = meta.block_metas[0].table_id();
682 let mut last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
683 for block_meta in &meta.block_metas {
684 let block_table_id = block_meta.table_id();
685 if last_table_id != block_table_id {
686 last_table_id = block_table_id;
687 last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
688 }
689
690 last_table_stats.total_compressed_size += block_meta.len as u64;
691 }
692 }
693
694 let writer_output = self
695 .writer
696 .finish(meta)
697 .instrument_await("sstable_writer_finish".verbose())
698 .await?;
699 let now = SystemTime::now()
708 .duration_since(SystemTime::UNIX_EPOCH)
709 .expect("Clock may have gone backwards")
710 .as_secs();
711 Ok(SstableBuilderOutput::<W::Output> {
712 sst_info: LocalSstableInfo::new(sst_info, self.table_stats, now),
713 writer_output,
714 stats: SstableBuilderOutputStats {
715 filter_size,
716 avg_key_size,
717 avg_value_size,
718 epoch_count: self.epoch_set.len(),
719 block_size_vec: self.block_size_vec,
720 sstable_file_size,
721 },
722 })
723 }
724
725 pub fn approximate_len(&self) -> usize {
726 self.writer.data_len()
727 + self.block_builder.approximate_len()
728 + self.filter_builder.approximate_len()
729 }
730
731 pub async fn build_block(&mut self) -> HummockResult<()> {
732 if self.block_builder.is_empty() {
734 return Ok(());
735 }
736
737 let block_meta = self.block_metas.last_mut().unwrap();
738 let uncompressed_block_size = self.block_builder.uncompressed_block_size();
739 block_meta.uncompressed_size = utils::checked_into_u32(uncompressed_block_size)
740 .unwrap_or_else(|_| {
741 panic!(
742 "WARN overflow can't convert uncompressed_block_size {} into u32 table {:?}",
743 uncompressed_block_size,
744 self.block_builder.table_id(),
745 )
746 });
747 let block = self.block_builder.build();
748 self.writer
749 .write_block(block, block_meta)
750 .instrument_await("sstable_write_block".verbose())
751 .await?;
752 self.block_size_vec.push(block.len());
753 let data_len = utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
754 panic!(
755 "WARN overflow can't convert writer_data_len {} into u32 table {:?}",
756 self.writer.data_len(),
757 self.block_builder.table_id(),
758 )
759 });
760 block_meta.len = data_len.checked_sub(block_meta.offset).unwrap_or_else(|| {
761 panic!(
762 "data_len should >= meta_offset, found data_len={}, meta_offset={}",
763 data_len, block_meta.offset
764 )
765 });
766
767 self.filter_builder
768 .switch_block(self.memory_limiter.clone());
769
770 if data_len as usize > self.options.capacity * 2 {
771 tracing::warn!(
772 "WARN unexpected block size {} table {:?}",
773 data_len,
774 self.block_builder.table_id()
775 );
776 }
777
778 self.block_builder.clear();
779 Ok(())
780 }
781
782 pub fn is_empty(&self) -> bool {
783 self.writer.data_len() > 0
784 }
785
786 pub fn reach_capacity(&self) -> bool {
788 self.approximate_len() >= self.options.capacity
789 }
790
791 fn finalize_last_table_stats(&mut self) {
792 if self.table_ids.is_empty() || self.last_table_id.is_none() {
793 return;
794 }
795 self.table_stats.insert(
796 self.last_table_id.unwrap(),
797 std::mem::take(&mut self.last_table_stats),
798 );
799 }
800}
801
802struct VnodeUserKeyRangeCollector {
804 max_bytes: usize,
805 current_size: usize,
806 ranges: BTreeMap<VirtualNode, (UserKey<Bytes>, UserKey<Bytes>)>,
807 current_vnode: VirtualNode,
808 range_start_key: Vec<u8>,
809}
810
811impl VnodeUserKeyRangeCollector {
812 fn new(max_bytes: usize) -> Self {
813 Self {
814 max_bytes,
815 current_size: 0,
816 ranges: BTreeMap::new(),
817 current_vnode: VirtualNode::ZERO,
818 range_start_key: Vec::new(),
819 }
820 }
821
822 fn with_limit(max_bytes: Option<usize>) -> Option<Self> {
823 max_bytes.filter(|&n| n > 0).map(Self::new)
824 }
825
826 fn observe_key(&mut self, vnode: VirtualNode, key: &[u8], prev_key: &[u8]) {
829 if self.current_size >= self.max_bytes {
830 return;
831 }
832
833 if self.range_start_key.is_empty() {
835 self.current_vnode = vnode;
836 self.range_start_key = key.to_vec();
837 return;
838 }
839
840 if vnode == self.current_vnode {
842 return;
843 }
844
845 self.seal_range(prev_key);
847
848 if self.current_size >= self.max_bytes {
850 return;
851 }
852
853 self.current_vnode = vnode;
855 self.range_start_key = key.to_vec();
856 }
857
858 fn seal_range(&mut self, right_key: &[u8]) {
860 let left_key = mem::take(&mut self.range_start_key);
861 self.current_size += mem::size_of::<VirtualNode>() + left_key.len() + right_key.len();
862
863 let left_full_key = FullKey::decode(&left_key);
864 let right_full_key = FullKey::decode(right_key);
865 let left_user_key = left_full_key.user_key.copy_into();
866 let right_user_key = right_full_key.user_key.copy_into();
867
868 assert_eq!(
873 left_user_key.get_vnode_id(),
874 right_user_key.get_vnode_id(),
875 "vnode changed within range: left_user {:?}, right_user {:?}",
876 left_user_key,
877 right_user_key
878 );
879 assert_eq!(
880 left_user_key.get_vnode_id(),
881 self.current_vnode.to_index(),
882 "vnode mismatch: left {:?}, right {:?}, expected vnode {}",
883 left_user_key,
884 right_user_key,
885 self.current_vnode.to_index()
886 );
887 assert_eq!(
888 left_user_key.table_id, right_user_key.table_id,
889 "table_id changed within range: left {:?}, right {:?}",
890 left_user_key, right_user_key
891 );
892
893 self.ranges
894 .insert(self.current_vnode, (left_user_key, right_user_key));
895 }
896
897 fn finish(mut self, last_key: &[u8]) -> Option<VnodeStatistics> {
899 if !self.range_start_key.is_empty() {
900 self.seal_range(last_key);
901 }
902
903 if self.ranges.len() > 1 {
904 let mut table_ids = self
906 .ranges
907 .values()
908 .flat_map(|(left, right)| [left.table_id, right.table_id]);
909 if let Some(first_table_id) = table_ids.next() {
910 for table_id in table_ids {
911 assert_eq!(
912 table_id, first_table_id,
913 "all vnode ranges must belong to the same table_id, found {:?} and {:?}",
914 table_id, first_table_id
915 );
916 }
917 }
918
919 Some(VnodeStatistics::from_map(self.ranges))
920 } else {
921 None
922 }
923 }
924}
925
926pub struct SstableBuilderOutputStats {
927 filter_size: usize,
928 avg_key_size: usize,
929 avg_value_size: usize,
930 epoch_count: usize,
931 block_size_vec: Vec<usize>, sstable_file_size: usize,
933}
934
935impl SstableBuilderOutputStats {
936 pub fn report_stats(&self, metrics: &Arc<CompactorMetrics>) {
937 if self.filter_size != 0 {
938 metrics
939 .sstable_bloom_filter_size
940 .observe(self.filter_size as _);
941 }
942
943 if self.sstable_file_size != 0 {
944 metrics
945 .sstable_file_size
946 .observe(self.sstable_file_size as _);
947 }
948
949 if self.avg_key_size != 0 {
950 metrics.sstable_avg_key_size.observe(self.avg_key_size as _);
951 }
952
953 if self.avg_value_size != 0 {
954 metrics
955 .sstable_avg_value_size
956 .observe(self.avg_value_size as _);
957 }
958
959 if self.epoch_count != 0 {
960 metrics
961 .sstable_distinct_epoch_count
962 .observe(self.epoch_count as _);
963 }
964
965 if !self.block_size_vec.is_empty() {
966 for block_size in &self.block_size_vec {
967 metrics.sstable_block_size.observe(*block_size as _);
968 }
969 }
970 }
971}
972
973#[cfg(test)]
974pub(super) mod tests {
975 use std::collections::{Bound, HashMap};
976
977 use risingwave_common::catalog::TableId;
978 use risingwave_common::hash::VirtualNode;
979 use risingwave_common::util::epoch::test_epoch;
980 use risingwave_hummock_sdk::key::UserKey;
981
982 use super::*;
983 use crate::assert_bytes_eq;
984 use crate::compaction_catalog_manager::{
985 CompactionCatalogAgent, DummyFilterKeyExtractor, MultiFilterKeyExtractor,
986 };
987 use crate::hummock::iterator::test_utils::mock_sstable_store;
988 use crate::hummock::sstable::xor_filter::BlockedXor16FilterBuilder;
989 use crate::hummock::test_utils::{
990 TEST_KEYS_COUNT, default_builder_opt_for_test, gen_test_sstable_impl, mock_sst_writer,
991 test_key_of, test_value_of,
992 };
993 use crate::hummock::{CachePolicy, Sstable, SstableWriterOptions, Xor8FilterBuilder};
994 use crate::monitor::StoreLocalStatistic;
995
996 #[tokio::test]
997 async fn test_empty() {
998 let opt = SstableBuilderOptions {
999 capacity: 0,
1000 block_capacity: 4096,
1001 restart_interval: 16,
1002 bloom_false_positive: 0.001,
1003 ..Default::default()
1004 };
1005
1006 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1007 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1008 let b = SstableBuilder::for_test(
1009 0,
1010 mock_sst_writer(&opt),
1011 opt,
1012 table_id_to_vnode,
1013 table_id_to_watermark_serde,
1014 );
1015
1016 b.finish().await.unwrap();
1017 }
1018
1019 fn encode_full_key(vnode: VirtualNode, table_key_suffix: &[u8]) -> Vec<u8> {
1020 let mut table_key = vnode.to_be_bytes().to_vec();
1021 table_key.extend_from_slice(table_key_suffix);
1022 FullKey::for_test(TableId::default(), table_key, 0).encode()
1023 }
1024
1025 fn table_key_of(vnode: VirtualNode, suffix: &[u8]) -> Vec<u8> {
1026 let mut key = vnode.to_be_bytes().to_vec();
1027 key.extend_from_slice(suffix);
1028 key
1029 }
1030
1031 #[test]
1032 fn test_vnode_user_key_range_basic_collection() {
1033 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1036 let vnode_1 = VirtualNode::from_index(1);
1037 let vnode_2 = VirtualNode::from_index(2);
1038
1039 let k1 = encode_full_key(vnode_1, b"k1");
1040 let k2 = encode_full_key(vnode_1, b"k2");
1041 let k3 = encode_full_key(vnode_2, b"k3");
1042 let k4 = encode_full_key(vnode_2, b"k4");
1043
1044 collector.observe_key(vnode_1, &k1, &[]);
1045 collector.observe_key(vnode_1, &k2, &k1);
1046 collector.observe_key(vnode_2, &k3, &k2);
1047 collector.observe_key(vnode_2, &k4, &k3);
1048
1049 let info = collector.finish(&k4).unwrap();
1050 assert_eq!(info.vnode_user_key_ranges().len(), 2);
1051
1052 let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1054 assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1055 assert_eq!(
1056 range1_right.table_key.as_ref(),
1057 table_key_of(vnode_1, b"k2")
1058 );
1059
1060 let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1062 assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1063 assert_eq!(
1064 range2_right.table_key.as_ref(),
1065 table_key_of(vnode_2, b"k4")
1066 );
1067 }
1068
1069 #[test]
1070 fn test_vnode_user_key_range_capacity_limit() {
1071 let vnode_1 = VirtualNode::from_index(1);
1081 let vnode_2 = VirtualNode::from_index(2);
1082 let vnode_3 = VirtualNode::from_index(3);
1083 let vnode_4 = VirtualNode::from_index(4);
1084
1085 let k1 = encode_full_key(vnode_1, b"k1");
1086 let k2 = encode_full_key(vnode_1, b"k2");
1087 let k3 = encode_full_key(vnode_2, b"k3");
1088 let k4 = encode_full_key(vnode_2, b"k4");
1089 let k5 = encode_full_key(vnode_3, b"k5");
1090 let k6 = encode_full_key(vnode_3, b"k6");
1091 let k7 = encode_full_key(vnode_4, b"k7");
1092
1093 let range_size = mem::size_of::<VirtualNode>() + k1.len() + k2.len();
1094 let estimated_next_size = mem::size_of::<VirtualNode>() + k3.len();
1095 let limit = range_size + estimated_next_size; let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(limit)).unwrap();
1097
1098 collector.observe_key(vnode_1, &k1, &[]);
1099 collector.observe_key(vnode_1, &k2, &k1);
1100 collector.observe_key(vnode_2, &k3, &k2); collector.observe_key(vnode_2, &k4, &k3);
1102 collector.observe_key(vnode_3, &k5, &k4); collector.observe_key(vnode_3, &k6, &k5); collector.observe_key(vnode_4, &k7, &k6); let info = collector.finish(&k7).unwrap();
1107 assert_eq!(
1108 info.vnode_user_key_ranges().len(),
1109 2,
1110 "Should collect exactly 2 vnodes"
1111 );
1112
1113 let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1115 assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1116 assert_eq!(
1117 range1_right.table_key.as_ref(),
1118 table_key_of(vnode_1, b"k2")
1119 );
1120
1121 let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1122 assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1123 assert_eq!(
1124 range2_right.table_key.as_ref(),
1125 table_key_of(vnode_2, b"k4")
1126 );
1127
1128 assert!(info.get_vnode_user_key_range(vnode_3).is_none());
1130 assert!(info.get_vnode_user_key_range(vnode_4).is_none());
1131 }
1132
1133 #[test]
1134 fn test_vnode_user_key_range_sparse_distribution() {
1135 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1138 let vnode_5 = VirtualNode::from_index(5);
1139 let vnode_10 = VirtualNode::from_index(10);
1140 let vnode_100 = VirtualNode::from_index(100);
1141
1142 let keys = vec![
1143 (vnode_5, encode_full_key(vnode_5, b"key_005_001")),
1144 (vnode_5, encode_full_key(vnode_5, b"key_005_002")),
1145 (vnode_5, encode_full_key(vnode_5, b"key_005_999")),
1146 (vnode_10, encode_full_key(vnode_10, b"key_010_001")),
1147 (vnode_10, encode_full_key(vnode_10, b"key_010_100")),
1148 (vnode_100, encode_full_key(vnode_100, b"key_100_001")),
1149 (vnode_100, encode_full_key(vnode_100, b"key_100_999")),
1150 ];
1151
1152 let mut prev_key = Vec::new();
1153 for (vnode, key) in &keys {
1154 collector.observe_key(*vnode, key, &prev_key);
1155 prev_key = key.clone();
1156 }
1157
1158 let info = collector.finish(&prev_key).unwrap();
1159 assert_eq!(info.vnode_user_key_ranges().len(), 3);
1160
1161 let (range5_left, range5_right) = info.get_vnode_user_key_range(vnode_5).unwrap();
1163 assert_eq!(
1164 range5_left.table_key.as_ref(),
1165 table_key_of(vnode_5, b"key_005_001")
1166 );
1167 assert_eq!(
1168 range5_right.table_key.as_ref(),
1169 table_key_of(vnode_5, b"key_005_999")
1170 );
1171
1172 let (range10_left, range10_right) = info.get_vnode_user_key_range(vnode_10).unwrap();
1173 assert_eq!(
1174 range10_left.table_key.as_ref(),
1175 table_key_of(vnode_10, b"key_010_001")
1176 );
1177 assert_eq!(
1178 range10_right.table_key.as_ref(),
1179 table_key_of(vnode_10, b"key_010_100")
1180 );
1181
1182 let (range100_left, range100_right) = info.get_vnode_user_key_range(vnode_100).unwrap();
1183 assert_eq!(
1184 range100_left.table_key.as_ref(),
1185 table_key_of(vnode_100, b"key_100_001")
1186 );
1187 assert_eq!(
1188 range100_right.table_key.as_ref(),
1189 table_key_of(vnode_100, b"key_100_999")
1190 );
1191
1192 assert!(
1194 info.get_vnode_user_key_range(VirtualNode::from_index(0))
1195 .is_none()
1196 );
1197 assert!(
1198 info.get_vnode_user_key_range(VirtualNode::from_index(7))
1199 .is_none()
1200 );
1201 assert!(
1202 info.get_vnode_user_key_range(VirtualNode::from_index(50))
1203 .is_none()
1204 );
1205 }
1206
1207 #[test]
1208 fn test_vnode_user_key_range_edge_cases() {
1209 assert!(VnodeUserKeyRangeCollector::with_limit(None).is_none());
1211 assert!(VnodeUserKeyRangeCollector::with_limit(Some(0)).is_none());
1212
1213 let collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1215 assert!(
1216 collector
1217 .finish(&encode_full_key(VirtualNode::ZERO, b"any"))
1218 .is_none()
1219 );
1220
1221 let vnode = VirtualNode::from_index(5);
1223 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1224 let a = encode_full_key(vnode, b"a");
1225 let b = encode_full_key(vnode, b"b");
1226 let c = encode_full_key(vnode, b"c");
1227 collector.observe_key(vnode, &a, &[]);
1228 collector.observe_key(vnode, &b, &a);
1229 collector.observe_key(vnode, &c, &b);
1230 assert!(
1231 collector.finish(&c).is_none(),
1232 "Single-vnode SST should not emit hints"
1233 );
1234
1235 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1)).unwrap();
1240 let vnode_1 = VirtualNode::from_index(1);
1241 let vnode_2 = VirtualNode::from_index(2);
1242 let key1 = encode_full_key(vnode_1, b"key1");
1243 let key2 = encode_full_key(vnode_1, b"key2");
1244 let key3 = encode_full_key(vnode_2, b"key3");
1245 collector.observe_key(vnode_1, &key1, &[]);
1246 collector.observe_key(vnode_1, &key2, &key1);
1247 collector.observe_key(vnode_2, &key3, &key2); assert!(
1249 collector.finish(&key3).is_none(),
1250 "Only 1 vnode collected, should return None"
1251 );
1252 }
1253
1254 #[tokio::test]
1255 async fn test_add_raw_block_compression_and_fallback() {
1256 let new_builder = |compression_algorithm| {
1257 let options = SstableBuilderOptions {
1258 block_capacity: 2048,
1259 compression_algorithm,
1260 ..Default::default()
1261 };
1262 SstableBuilder::<_, BlockedXor16FilterBuilder>::new(
1263 1,
1264 mock_sst_writer(&options),
1265 BlockedXor16FilterBuilder::create(options.filter_builder_options()),
1266 options,
1267 CompactionCatalogAgent::for_test(vec![0]),
1268 None,
1269 )
1270 };
1271 let value = vec![7; 600];
1272 let mut source = new_builder(CompressionAlgorithm::None);
1273 for index in [2, 3] {
1274 source
1275 .add(test_key_of(index).to_ref(), HummockValue::put(&value))
1276 .await
1277 .unwrap();
1278 }
1279 let (data, meta) = source.finish().await.unwrap().writer_output;
1280 let source = Sstable::new(1.into(), meta, false);
1281 assert_eq!(source.meta.block_metas.len(), 1);
1282 let source_meta = &source.meta.block_metas[0];
1283 let raw = data.slice(..source_meta.len as usize);
1284
1285 for compression in [CompressionAlgorithm::Lz4, CompressionAlgorithm::Zstd] {
1286 for prefix_size in [16, 600] {
1287 let mut builder = new_builder(compression);
1288 let prefix_value = vec![1; prefix_size];
1289 builder
1290 .add(test_key_of(1).to_ref(), HummockValue::put(&prefix_value))
1291 .await
1292 .unwrap();
1293 let written = builder
1294 .add_raw_block(
1295 raw.clone(),
1296 source.filter_reader.get_block_raw_filter(0),
1297 test_key_of(2),
1298 test_key_of(3).encode(),
1299 source_meta.clone(),
1300 )
1301 .await
1302 .unwrap();
1303 let (data, meta) = builder.finish().await.unwrap().writer_output;
1304 if prefix_size == 16 {
1305 assert_eq!(written, None);
1307 assert_eq!(meta.block_metas.len(), 1);
1308 } else {
1309 assert_eq!(meta.block_metas.len(), 2);
1310 let copied = &meta.block_metas[1];
1311 assert_eq!(written, Some(copied.len as usize));
1312 assert!(copied.len < source_meta.len);
1313 assert_eq!(copied.uncompressed_size, source_meta.uncompressed_size);
1314 }
1315 let mut expected_index = 1;
1316 let mut offset = 0;
1317 for block_meta in &meta.block_metas {
1318 assert_eq!(block_meta.offset, offset);
1319 offset += block_meta.len;
1320 let block = data.slice(block_meta.offset as usize..offset as usize);
1321 assert_eq!(Block::get_algorithm(&block).unwrap(), compression);
1322 let block =
1323 Block::decode(block, block_meta.uncompressed_size as usize).unwrap();
1324 let mut iter =
1325 BlockIterator::new(BlockHolder::from_owned_block(Box::new(block)));
1326 iter.seek_to_first();
1327 while iter.is_valid() {
1328 assert_eq!(iter.key(), test_key_of(expected_index).to_ref());
1329 let expected_value = if expected_index == 1 {
1330 &prefix_value
1331 } else {
1332 &value
1333 };
1334 assert_eq!(
1335 HummockValue::from_slice(iter.value()).unwrap(),
1336 HummockValue::put(expected_value.as_slice())
1337 );
1338 expected_index += 1;
1339 iter.next();
1340 }
1341 }
1342 assert_eq!(expected_index, 4);
1343 }
1344 }
1345 }
1346
1347 #[tokio::test]
1348 async fn test_basic() {
1349 let opt = default_builder_opt_for_test();
1350
1351 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1352 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1353 let mut b = SstableBuilder::for_test(
1354 0,
1355 mock_sst_writer(&opt),
1356 opt,
1357 table_id_to_vnode,
1358 table_id_to_watermark_serde,
1359 );
1360
1361 for i in 0..TEST_KEYS_COUNT {
1362 b.add_for_test(
1363 test_key_of(i).to_ref(),
1364 HummockValue::put(&test_value_of(i)),
1365 )
1366 .await
1367 .unwrap();
1368 }
1369
1370 let output = b.finish().await.unwrap();
1371 let info = output.sst_info.sst_info;
1372
1373 assert_bytes_eq!(test_key_of(0).encode(), info.key_range.left);
1374 assert_bytes_eq!(
1375 test_key_of(TEST_KEYS_COUNT - 1).encode(),
1376 info.key_range.right
1377 );
1378 let (data, meta) = output.writer_output;
1379 assert_eq!(info.file_size, meta.estimated_size as u64);
1380 let offset = info.meta_offset as usize;
1381 let meta2 = SstableMeta::decode(&data[offset..]).unwrap();
1382 assert_eq!(meta2, meta);
1383 }
1384
1385 async fn test_with_xor_filter_builder<F: FilterBuilder>(
1386 bloom_false_positive: f64,
1387 expected_filter_type: PbSstableFilterType,
1388 expected_filter_layout: PbSstableFilterLayout,
1389 ) {
1390 let key_count = 1000;
1391
1392 let opts = SstableBuilderOptions {
1393 capacity: 0,
1394 block_capacity: 4096,
1395 restart_interval: 16,
1396 bloom_false_positive,
1397 ..Default::default()
1398 };
1399
1400 let sstable_store = mock_sstable_store().await;
1402 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1403 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1404 let sst_info = gen_test_sstable_impl::<Vec<u8>, F>(
1405 opts,
1406 0,
1407 (0..TEST_KEYS_COUNT).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
1408 sstable_store.clone(),
1409 CachePolicy::NotFill,
1410 table_id_to_vnode,
1411 table_id_to_watermark_serde,
1412 )
1413 .await;
1414 assert_eq!(sst_info.filter_type, expected_filter_type);
1415 assert_eq!(sst_info.filter_layout, expected_filter_layout);
1416 let table = sstable_store
1417 .sstable(&sst_info, &mut StoreLocalStatistic::default())
1418 .await
1419 .unwrap();
1420
1421 assert!(table.has_filter());
1422 for i in 0..key_count {
1423 let full_key = test_key_of(i);
1424 let hash = Sstable::hash_for_filter(full_key.user_key.encode().as_slice(), 0);
1425 let key_ref = full_key.user_key.as_ref();
1426 assert!(
1427 table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash),
1428 "failed at {}",
1429 i
1430 );
1431 }
1432 }
1433
1434 #[tokio::test]
1435 async fn test_xor_filter_builder_output() {
1436 test_with_xor_filter_builder::<Xor16FilterBuilder>(
1437 0.0,
1438 PbSstableFilterType::SstableFilterXor16,
1439 PbSstableFilterLayout::Plain,
1440 )
1441 .await;
1442 test_with_xor_filter_builder::<Xor16FilterBuilder>(
1443 0.01,
1444 PbSstableFilterType::SstableFilterXor16,
1445 PbSstableFilterLayout::Plain,
1446 )
1447 .await;
1448 test_with_xor_filter_builder::<Xor8FilterBuilder>(
1449 0.01,
1450 PbSstableFilterType::SstableFilterXor8,
1451 PbSstableFilterLayout::Plain,
1452 )
1453 .await;
1454 test_with_xor_filter_builder::<BlockedXor16FilterBuilder>(
1455 0.01,
1456 PbSstableFilterType::SstableFilterXor16,
1457 PbSstableFilterLayout::Blocked,
1458 )
1459 .await;
1460 }
1461
1462 #[tokio::test]
1463 async fn test_no_xor_filter_block() {
1464 let opts = SstableBuilderOptions::default();
1465 let sstable_store = mock_sstable_store().await;
1467 let writer_opts = SstableWriterOptions::default();
1468 let object_id = 1;
1469 let writer = sstable_store
1470 .clone()
1471 .create_sst_writer(object_id, writer_opts);
1472 let mut filter = MultiFilterKeyExtractor::default();
1473 filter.register(
1474 1.into(),
1475 FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1476 );
1477 filter.register(
1478 2.into(),
1479 FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
1480 );
1481 filter.register(
1482 3.into(),
1483 FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1484 );
1485
1486 let table_id_to_vnode = HashMap::from_iter(vec![
1487 (1.into(), VirtualNode::COUNT_FOR_TEST),
1488 (2.into(), VirtualNode::COUNT_FOR_TEST),
1489 (3.into(), VirtualNode::COUNT_FOR_TEST),
1490 ]);
1491 let table_id_to_watermark_serde =
1492 HashMap::from_iter(vec![(1.into(), None), (2.into(), None), (3.into(), None)]);
1493
1494 let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
1495 FilterKeyExtractorImpl::Multi(filter),
1496 table_id_to_vnode,
1497 table_id_to_watermark_serde,
1498 HashMap::default(),
1499 ));
1500
1501 let mut builder = SstableBuilder::new(
1502 object_id,
1503 writer,
1504 BlockedXor16FilterBuilder::create(opts.filter_builder_options()),
1505 opts,
1506 compaction_catalog_agent_ref,
1507 None,
1508 );
1509
1510 let key_count: usize = 10000;
1511 for table_id in 1..4 {
1512 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1513 for idx in 0..key_count {
1514 table_key.resize(VirtualNode::SIZE, 0);
1515 table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1516 let k = UserKey::for_test(TableId::new(table_id), table_key.as_ref());
1517 let v = test_value_of(idx);
1518 builder
1519 .add(
1520 FullKey::from_user_key(k, test_epoch(1)),
1521 HummockValue::put(v.as_ref()),
1522 )
1523 .await
1524 .unwrap();
1525 }
1526 }
1527 let ret = builder.finish().await.unwrap();
1528 let sst_info = ret.sst_info.sst_info.clone();
1529 ret.writer_output.await.unwrap().unwrap();
1530 let table = sstable_store
1531 .sstable(&sst_info, &mut StoreLocalStatistic::default())
1532 .await
1533 .unwrap();
1534 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1535 for idx in 0..key_count {
1536 table_key.resize(VirtualNode::SIZE, 0);
1537 table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1538 let k = UserKey::for_test(TableId::new(2), table_key.as_slice());
1539 let hash = Sstable::hash_for_filter(&k.encode(), 2);
1540 let key_ref = k.as_ref();
1541 assert!(
1542 table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash)
1543 );
1544 }
1545 }
1546}