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: BytesMut,
153 raw_value: BytesMut,
154 last_table_id: Option<TableId>,
155 sst_object_id: HummockSstableObjectId,
156
157 table_stats: TableStatsMap,
159 last_table_stats: TableStats,
162
163 filter_builder: F,
164
165 epoch_set: BTreeSet<u64>,
166 memory_limiter: Option<Arc<MemoryLimiter>>,
167
168 block_size_vec: Vec<usize>, vnode_range_collector: Option<VnodeUserKeyRangeCollector>,
170}
171
172impl<W: SstableWriter> SstableBuilder<W, Xor16FilterBuilder> {
173 pub fn for_test(
174 sstable_id: u64,
175 writer: W,
176 options: SstableBuilderOptions,
177 table_id_to_vnode: HashMap<impl Into<TableId>, usize>,
178 table_id_to_watermark_serde: HashMap<
179 impl Into<TableId>,
180 Option<(OrderedRowSerde, OrderedRowSerde, usize)>,
181 >,
182 ) -> Self {
183 let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
184 FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
185 table_id_to_vnode
186 .into_iter()
187 .map(|(table_id, v)| (table_id.into(), v))
188 .collect(),
189 table_id_to_watermark_serde
190 .into_iter()
191 .map(|(table_id, v)| (table_id.into(), v))
192 .collect(),
193 HashMap::default(),
194 ));
195
196 Self::new(
197 sstable_id,
198 writer,
199 Xor16FilterBuilder::create(options.filter_builder_options()),
200 options,
201 compaction_catalog_agent_ref,
202 None,
203 )
204 }
205}
206
207impl<W: SstableWriter, F: FilterBuilder> SstableBuilder<W, F> {
208 pub fn new(
209 sst_object_id: impl Into<HummockSstableObjectId>,
210 writer: W,
211 filter_builder: F,
212 options: SstableBuilderOptions,
213 compaction_catalog_agent_ref: CompactionCatalogAgentRef,
214 memory_limiter: Option<Arc<MemoryLimiter>>,
215 ) -> Self {
216 let sst_object_id = sst_object_id.into();
217 Self {
218 vnode_range_collector: VnodeUserKeyRangeCollector::with_limit(
219 options.max_vnode_key_range_bytes,
220 ),
221 options: options.clone(),
222 writer,
223 block_builder: BlockBuilder::new(BlockBuilderOptions {
224 capacity: options.block_capacity,
225 restart_interval: options.restart_interval,
226 compression_algorithm: options.compression_algorithm,
227 }),
228 filter_builder,
229 block_metas: Vec::with_capacity(options.capacity / options.block_capacity + 1),
230 table_ids: BTreeSet::new(),
231 last_table_id: None,
232 raw_key: BytesMut::new(),
233 raw_value: BytesMut::new(),
234 last_full_key: vec![],
235 sst_object_id,
236 compaction_catalog_agent_ref,
237 table_stats: Default::default(),
238 last_table_stats: Default::default(),
239 epoch_set: BTreeSet::default(),
240 memory_limiter,
241 block_size_vec: Vec::new(),
242 }
243 }
244
245 pub async fn add_for_test(
247 &mut self,
248 full_key: FullKey<&[u8]>,
249 value: HummockValue<&[u8]>,
250 ) -> HummockResult<()> {
251 self.add(full_key, value).await
252 }
253
254 pub fn current_block_size(&self) -> usize {
255 self.block_builder.approximate_len()
256 }
257
258 pub async fn add_raw_block(
260 &mut self,
261 buf: Bytes,
262 filter_data: Vec<u8>,
263 smallest_key: FullKey<Vec<u8>>,
264 largest_key: Vec<u8>,
265 mut meta: BlockMeta,
266 ) -> HummockResult<bool> {
267 let table_id = smallest_key.user_key.table_id;
268 if self.last_table_id.is_none() || self.last_table_id.unwrap() != table_id {
269 if !self.block_builder.is_empty() {
270 self.build_block()
272 .instrument_await("sstable_build_block_on_raw_table_switch".verbose())
273 .await?;
274 }
275
276 self.table_ids.insert(table_id);
277 self.finalize_last_table_stats();
278 self.last_table_id = Some(table_id);
279 }
280
281 if !self.block_builder.is_empty() {
282 let min_block_size = std::cmp::min(MIN_BLOCK_SIZE, self.options.block_capacity / 4);
283
284 if self.block_builder.approximate_len() < min_block_size {
286 let block = Block::decode(buf, meta.uncompressed_size as usize)?;
287 let mut iter = BlockIterator::new(BlockHolder::from_owned_block(Box::new(block)));
288 iter.seek_to_first();
289 while iter.is_valid() {
290 let value = HummockValue::from_slice(iter.value()).unwrap_or_else(|_| {
291 panic!(
292 "decode failed for fast compact sst_id {} block_idx {} last_table_id {:?}",
293 self.sst_object_id, self.block_metas.len(), self.last_table_id
294 )
295 });
296 self.add_impl(iter.key(), value, false)
297 .instrument_await("sstable_add_raw_block_fallback_add".verbose())
298 .await?;
299 iter.next();
300 }
301 return Ok(false);
302 }
303
304 self.build_block()
305 .instrument_await("sstable_build_block_before_raw_block".verbose())
306 .await?;
307 }
308 self.last_full_key = largest_key;
309 assert_eq!(
310 meta.len as usize,
311 buf.len(),
312 "meta {} buf {} last_table_id {:?}",
313 meta.len,
314 buf.len(),
315 self.last_table_id
316 );
317 meta.offset = self.writer.data_len() as u32;
318 self.block_metas.push(meta);
319 self.filter_builder.add_raw_data(filter_data);
320 let block_meta = self.block_metas.last_mut().unwrap();
321 self.writer
322 .write_block_bytes(buf, block_meta)
323 .instrument_await("sstable_write_raw_block_bytes".verbose())
324 .await?;
325
326 Ok(true)
327 }
328
329 pub async fn add(
331 &mut self,
332 full_key: FullKey<&[u8]>,
333 value: HummockValue<&[u8]>,
334 ) -> HummockResult<()> {
335 self.add_impl(full_key, value, true).await
336 }
337
338 async fn add_impl(
340 &mut self,
341 full_key: FullKey<&[u8]>,
342 value: HummockValue<&[u8]>,
343 could_switch_block: bool,
344 ) -> HummockResult<()> {
345 const LARGE_KEY_LEN: usize = MAX_KEY_LEN >> 1;
346
347 let table_key_len = full_key.user_key.table_key.as_ref().len();
348 let table_value_len = match &value {
349 HummockValue::Put(t) => t.len(),
350 HummockValue::Delete => 0,
351 };
352 let large_value_len = self.options.max_sst_size as usize / 10;
353 let large_key_value_len = self.options.max_sst_size as usize / 2;
354 if table_key_len >= LARGE_KEY_LEN
355 || table_value_len > large_value_len
356 || table_key_len + table_value_len > large_key_value_len
357 {
358 let table_id = full_key.user_key.table_id;
359 tracing::warn!(
360 "A large key/value (table_id={}, key len={}, value len={}, epoch={}, spill offset={}) is added to block",
361 table_id,
362 table_key_len,
363 table_value_len,
364 full_key.epoch_with_gap.pure_epoch(),
365 full_key.epoch_with_gap.offset(),
366 );
367 }
368
369 full_key.encode_into(&mut self.raw_key);
371 value.encode(&mut self.raw_value);
372 let is_new_user_key = self.last_full_key.is_empty()
373 || !user_key(&self.raw_key).eq(user_key(self.last_full_key.as_slice()));
374 let table_id = full_key.user_key.table_id;
375 let is_new_table = self.last_table_id != Some(table_id);
376 let current_block_size = self.current_block_size();
377 let is_block_full = current_block_size >= self.options.block_capacity
378 || (current_block_size > self.options.block_capacity / 4 * 3
379 && current_block_size + self.raw_value.len() + self.raw_key.len()
380 > self.options.block_capacity);
381
382 if is_new_table {
383 assert!(
384 could_switch_block,
385 "is_new_user_key {} sst_id {} block_idx {} table_id {} last_table_id {:?} full_key {:?}",
386 is_new_user_key,
387 self.sst_object_id,
388 self.block_metas.len(),
389 table_id,
390 self.last_table_id,
391 full_key
392 );
393 self.table_ids.insert(table_id);
394 self.finalize_last_table_stats();
395 self.last_table_id = Some(table_id);
396 if !self.block_builder.is_empty() {
397 self.build_block()
398 .instrument_await("sstable_build_block_on_table_switch".verbose())
399 .await?;
400 }
401 } else if is_block_full && could_switch_block {
402 self.build_block()
403 .instrument_await("sstable_build_block_on_block_full".verbose())
404 .await?;
405 }
406 self.last_table_stats.total_key_count += 1;
407 self.epoch_set.insert(full_key.epoch_with_gap.pure_epoch());
408
409 if self.block_builder.is_empty() {
411 let smallest_key = if let Some(threshold) =
412 self.options.shorten_block_meta_key_threshold
413 && !self.last_full_key.is_empty()
414 && full_key.encoded_len() >= threshold
415 {
416 let prev = FullKey::decode(&self.last_full_key);
417 if let Some(shortened) = try_shorten_block_smallest_key(&prev, &full_key) {
418 shortened.encode()
419 } else {
420 full_key.encode()
421 }
422 } else {
423 full_key.encode()
424 };
425
426 self.block_metas.push(BlockMeta {
427 offset: utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
428 panic!(
429 "WARN overflow can't convert writer_data_len {} into u32 sst_id {} block_idx {} tables {:?}",
430 self.writer.data_len(),
431 self.sst_object_id,
432 self.block_metas.len(),
433 self.table_ids,
434 )
435 }),
436 len: 0,
437 smallest_key,
438 uncompressed_size: 0,
439 total_key_count: 0,
440 stale_key_count: 0,
441 });
442 }
443
444 let filter_key = self
445 .compaction_catalog_agent_ref
446 .extract(user_key(&self.raw_key));
447 if !filter_key.is_empty() {
449 self.filter_builder
450 .add_key(filter_key, table_id.as_raw_id());
451 }
452 self.block_builder.add(
454 table_id,
455 &self.raw_key[TABLE_PREFIX_LEN..],
456 self.raw_value.as_ref(),
457 );
458 self.block_metas.last_mut().unwrap().total_key_count += 1;
459 if !is_new_user_key || value.is_delete() {
460 self.block_metas.last_mut().unwrap().stale_key_count += 1;
461 }
462 self.last_table_stats.total_key_size += full_key.encoded_len() as i64;
463 self.last_table_stats.total_value_size += value.encoded_len() as i64;
464
465 if let Some(collector) = self.vnode_range_collector.as_mut() {
466 collector.observe_key(
467 VirtualNode::from_index(full_key.user_key.get_vnode_id()),
468 &self.raw_key,
469 &self.last_full_key,
470 );
471 }
472
473 self.last_full_key.clear();
474 self.last_full_key.extend_from_slice(&self.raw_key);
475
476 self.raw_key.clear();
477 self.raw_value.clear();
478 Ok(())
479 }
480
481 pub async fn finish(mut self) -> HummockResult<SstableBuilderOutput<W::Output>> {
491 let smallest_key = if self.block_metas.is_empty() {
492 vec![]
493 } else {
494 self.block_metas[0].smallest_key.clone()
495 };
496 let largest_key = self.last_full_key.clone();
497 self.finalize_last_table_stats();
498
499 assert!(
502 self.table_ids.len() <= 1 || self.vnode_range_collector.is_none(),
503 "vnode key-range hints are only supported for single-table SSTs, found {} tables",
504 self.table_ids.len()
505 );
506
507 self.build_block()
508 .instrument_await("sstable_finish_build_block".verbose())
509 .await?;
510 let right_exclusive = false;
511 let meta_offset = self.writer.data_len() as u64;
512
513 let filter_data = self.filter_builder.finish(self.memory_limiter.clone());
514 let (filter_type, filter_layout) = if filter_data.is_none() {
515 (
516 PbSstableFilterType::SstableFilterNone,
517 PbSstableFilterLayout::Unspecified,
518 )
519 } else if self.filter_builder.support_blocked_raw_data() {
520 (
521 self.filter_builder.filter_type(),
522 PbSstableFilterLayout::Blocked,
523 )
524 } else {
525 (
526 self.filter_builder.filter_type(),
527 PbSstableFilterLayout::Plain,
528 )
529 };
530
531 let total_key_count = self
532 .block_metas
533 .iter()
534 .map(|block_meta| block_meta.total_key_count as u64)
535 .sum::<u64>();
536 let stale_key_count = self
537 .block_metas
538 .iter()
539 .map(|block_meta| block_meta.stale_key_count as u64)
540 .sum::<u64>();
541 let uncompressed_file_size = self
542 .block_metas
543 .iter()
544 .map(|block_meta| block_meta.uncompressed_size as u64)
545 .sum::<u64>();
546
547 #[expect(deprecated)]
548 let mut meta = SstableMeta {
549 block_metas: self.block_metas,
550 bloom_filter: filter_data.unwrap_or_default(),
551 estimated_size: 0,
552 key_count: utils::checked_into_u32(total_key_count).unwrap_or_else(|_| {
553 panic!(
554 "WARN overflow can't convert total_key_count {} into u32 tables {:?}",
555 total_key_count, self.table_ids,
556 )
557 }),
558 smallest_key,
559 largest_key,
560 version: VERSION,
561 meta_offset,
562 monotonic_tombstone_events: vec![],
563 };
564
565 let meta_encode_size = meta.encoded_size();
566 let encoded_size_u32 = utils::checked_into_u32(meta_encode_size).unwrap_or_else(|_| {
567 panic!(
568 "WARN overflow can't convert meta_encoded_size {} into u32 tables {:?}",
569 meta_encode_size, self.table_ids,
570 )
571 });
572 let meta_offset_u32 = utils::checked_into_u32(meta_offset).unwrap_or_else(|_| {
573 panic!(
574 "WARN overflow can't convert meta_offset {} into u32 tables {:?}",
575 meta_offset, self.table_ids,
576 )
577 });
578 meta.estimated_size = encoded_size_u32
579 .checked_add(meta_offset_u32)
580 .unwrap_or_else(|| {
581 panic!(
582 "WARN overflow encoded_size_u32 {} meta_offset_u32 {} table_id {:?} table_ids {:?}",
583 encoded_size_u32, meta_offset_u32, self.last_table_id, self.table_ids
584 )
585 });
586
587 let (avg_key_size, avg_value_size) = if self.table_stats.is_empty() {
588 (0, 0)
589 } else {
590 let total_key_count: usize = self
591 .table_stats
592 .values()
593 .map(|s| s.total_key_count as usize)
594 .sum();
595
596 let total_key_size: usize = self
597 .table_stats
598 .values()
599 .map(|s| s.total_key_size as usize)
600 .sum();
601
602 let total_value_size: usize = self
603 .table_stats
604 .values()
605 .map(|s| s.total_value_size as usize)
606 .sum();
607
608 (
609 total_key_size.checked_div(total_key_count).unwrap_or(0),
610 total_value_size.checked_div(total_key_count).unwrap_or(0),
611 )
612 };
613
614 let (min_epoch, max_epoch) = {
615 if self.epoch_set.is_empty() {
616 (HummockEpoch::MAX, u64::MIN)
617 } else {
618 (
619 *self.epoch_set.first().unwrap(),
620 *self.epoch_set.last().unwrap(),
621 )
622 }
623 };
624
625 let vnode_user_key_ranges = self
626 .vnode_range_collector
627 .take()
628 .and_then(|collector| collector.finish(&self.last_full_key));
629
630 let sst_info: SstableInfo = SstableInfoInner {
631 object_id: self.sst_object_id,
632 sst_id: self.sst_object_id.as_raw_id().into(),
634 key_range: KeyRange {
635 left: Bytes::from(meta.smallest_key.clone()),
636 right: Bytes::from(meta.largest_key.clone()),
637 right_exclusive,
638 },
639 file_size: meta.estimated_size as u64,
640 table_ids: self.table_ids.into_iter().collect(),
641 meta_offset: meta.meta_offset,
642 stale_key_count,
643 total_key_count,
644 uncompressed_file_size: uncompressed_file_size + meta.encoded_size() as u64,
645 min_epoch,
646 max_epoch,
647 range_tombstone_count: 0,
648 sst_size: meta.estimated_size as u64,
649 filter_type,
650 filter_layout,
651 vnode_statistics: vnode_user_key_ranges,
652 }
653 .into();
654
655 tracing::trace!(
656 "meta_size {} filter_size {} add_key_counts {} stale_key_count {} min_epoch {} max_epoch {} epoch_count {}",
657 meta.encoded_size(),
658 meta.bloom_filter.len(),
659 total_key_count,
660 stale_key_count,
661 min_epoch,
662 max_epoch,
663 self.epoch_set.len()
664 );
665 let filter_size = meta.bloom_filter.len();
666 let sstable_file_size = sst_info.file_size as usize;
667
668 if !meta.block_metas.is_empty() {
669 let mut last_table_id = meta.block_metas[0].table_id();
671 let mut last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
672 for block_meta in &meta.block_metas {
673 let block_table_id = block_meta.table_id();
674 if last_table_id != block_table_id {
675 last_table_id = block_table_id;
676 last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
677 }
678
679 last_table_stats.total_compressed_size += block_meta.len as u64;
680 }
681 }
682
683 let writer_output = self
684 .writer
685 .finish(meta)
686 .instrument_await("sstable_writer_finish".verbose())
687 .await?;
688 let now = SystemTime::now()
697 .duration_since(SystemTime::UNIX_EPOCH)
698 .expect("Clock may have gone backwards")
699 .as_secs();
700 Ok(SstableBuilderOutput::<W::Output> {
701 sst_info: LocalSstableInfo::new(sst_info, self.table_stats, now),
702 writer_output,
703 stats: SstableBuilderOutputStats {
704 filter_size,
705 avg_key_size,
706 avg_value_size,
707 epoch_count: self.epoch_set.len(),
708 block_size_vec: self.block_size_vec,
709 sstable_file_size,
710 },
711 })
712 }
713
714 pub fn approximate_len(&self) -> usize {
715 self.writer.data_len()
716 + self.block_builder.approximate_len()
717 + self.filter_builder.approximate_len()
718 }
719
720 pub async fn build_block(&mut self) -> HummockResult<()> {
721 if self.block_builder.is_empty() {
723 return Ok(());
724 }
725
726 let block_meta = self.block_metas.last_mut().unwrap();
727 let uncompressed_block_size = self.block_builder.uncompressed_block_size();
728 block_meta.uncompressed_size = utils::checked_into_u32(uncompressed_block_size)
729 .unwrap_or_else(|_| {
730 panic!(
731 "WARN overflow can't convert uncompressed_block_size {} into u32 table {:?}",
732 uncompressed_block_size,
733 self.block_builder.table_id(),
734 )
735 });
736 let block = self.block_builder.build();
737 self.writer
738 .write_block(block, block_meta)
739 .instrument_await("sstable_write_block".verbose())
740 .await?;
741 self.block_size_vec.push(block.len());
742 let data_len = utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
743 panic!(
744 "WARN overflow can't convert writer_data_len {} into u32 table {:?}",
745 self.writer.data_len(),
746 self.block_builder.table_id(),
747 )
748 });
749 block_meta.len = data_len.checked_sub(block_meta.offset).unwrap_or_else(|| {
750 panic!(
751 "data_len should >= meta_offset, found data_len={}, meta_offset={}",
752 data_len, block_meta.offset
753 )
754 });
755
756 self.filter_builder
757 .switch_block(self.memory_limiter.clone());
758
759 if data_len as usize > self.options.capacity * 2 {
760 tracing::warn!(
761 "WARN unexpected block size {} table {:?}",
762 data_len,
763 self.block_builder.table_id()
764 );
765 }
766
767 self.block_builder.clear();
768 Ok(())
769 }
770
771 pub fn is_empty(&self) -> bool {
772 self.writer.data_len() > 0
773 }
774
775 pub fn reach_capacity(&self) -> bool {
777 self.approximate_len() >= self.options.capacity
778 }
779
780 fn finalize_last_table_stats(&mut self) {
781 if self.table_ids.is_empty() || self.last_table_id.is_none() {
782 return;
783 }
784 self.table_stats.insert(
785 self.last_table_id.unwrap(),
786 std::mem::take(&mut self.last_table_stats),
787 );
788 }
789}
790
791struct VnodeUserKeyRangeCollector {
793 max_bytes: usize,
794 current_size: usize,
795 ranges: BTreeMap<VirtualNode, (UserKey<Bytes>, UserKey<Bytes>)>,
796 current_vnode: VirtualNode,
797 range_start_key: Vec<u8>,
798}
799
800impl VnodeUserKeyRangeCollector {
801 fn new(max_bytes: usize) -> Self {
802 Self {
803 max_bytes,
804 current_size: 0,
805 ranges: BTreeMap::new(),
806 current_vnode: VirtualNode::ZERO,
807 range_start_key: Vec::new(),
808 }
809 }
810
811 fn with_limit(max_bytes: Option<usize>) -> Option<Self> {
812 max_bytes.filter(|&n| n > 0).map(Self::new)
813 }
814
815 fn observe_key(&mut self, vnode: VirtualNode, key: &[u8], prev_key: &[u8]) {
818 if self.current_size >= self.max_bytes {
819 return;
820 }
821
822 if self.range_start_key.is_empty() {
824 self.current_vnode = vnode;
825 self.range_start_key = key.to_vec();
826 return;
827 }
828
829 if vnode == self.current_vnode {
831 return;
832 }
833
834 self.seal_range(prev_key);
836
837 if self.current_size >= self.max_bytes {
839 return;
840 }
841
842 self.current_vnode = vnode;
844 self.range_start_key = key.to_vec();
845 }
846
847 fn seal_range(&mut self, right_key: &[u8]) {
849 let left_key = mem::take(&mut self.range_start_key);
850 self.current_size += mem::size_of::<VirtualNode>() + left_key.len() + right_key.len();
851
852 let left_full_key = FullKey::decode(&left_key);
853 let right_full_key = FullKey::decode(right_key);
854 let left_user_key = left_full_key.user_key.copy_into();
855 let right_user_key = right_full_key.user_key.copy_into();
856
857 assert_eq!(
862 left_user_key.get_vnode_id(),
863 right_user_key.get_vnode_id(),
864 "vnode changed within range: left_user {:?}, right_user {:?}",
865 left_user_key,
866 right_user_key
867 );
868 assert_eq!(
869 left_user_key.get_vnode_id(),
870 self.current_vnode.to_index(),
871 "vnode mismatch: left {:?}, right {:?}, expected vnode {}",
872 left_user_key,
873 right_user_key,
874 self.current_vnode.to_index()
875 );
876 assert_eq!(
877 left_user_key.table_id, right_user_key.table_id,
878 "table_id changed within range: left {:?}, right {:?}",
879 left_user_key, right_user_key
880 );
881
882 self.ranges
883 .insert(self.current_vnode, (left_user_key, right_user_key));
884 }
885
886 fn finish(mut self, last_key: &[u8]) -> Option<VnodeStatistics> {
888 if !self.range_start_key.is_empty() {
889 self.seal_range(last_key);
890 }
891
892 if self.ranges.len() > 1 {
893 let mut table_ids = self
895 .ranges
896 .values()
897 .flat_map(|(left, right)| [left.table_id, right.table_id]);
898 if let Some(first_table_id) = table_ids.next() {
899 for table_id in table_ids {
900 assert_eq!(
901 table_id, first_table_id,
902 "all vnode ranges must belong to the same table_id, found {:?} and {:?}",
903 table_id, first_table_id
904 );
905 }
906 }
907
908 Some(VnodeStatistics::from_map(self.ranges))
909 } else {
910 None
911 }
912 }
913}
914
915pub struct SstableBuilderOutputStats {
916 filter_size: usize,
917 avg_key_size: usize,
918 avg_value_size: usize,
919 epoch_count: usize,
920 block_size_vec: Vec<usize>, sstable_file_size: usize,
922}
923
924impl SstableBuilderOutputStats {
925 pub fn report_stats(&self, metrics: &Arc<CompactorMetrics>) {
926 if self.filter_size != 0 {
927 metrics
928 .sstable_bloom_filter_size
929 .observe(self.filter_size as _);
930 }
931
932 if self.sstable_file_size != 0 {
933 metrics
934 .sstable_file_size
935 .observe(self.sstable_file_size as _);
936 }
937
938 if self.avg_key_size != 0 {
939 metrics.sstable_avg_key_size.observe(self.avg_key_size as _);
940 }
941
942 if self.avg_value_size != 0 {
943 metrics
944 .sstable_avg_value_size
945 .observe(self.avg_value_size as _);
946 }
947
948 if self.epoch_count != 0 {
949 metrics
950 .sstable_distinct_epoch_count
951 .observe(self.epoch_count as _);
952 }
953
954 if !self.block_size_vec.is_empty() {
955 for block_size in &self.block_size_vec {
956 metrics.sstable_block_size.observe(*block_size as _);
957 }
958 }
959 }
960}
961
962#[cfg(test)]
963pub(super) mod tests {
964 use std::collections::{Bound, HashMap};
965
966 use risingwave_common::catalog::TableId;
967 use risingwave_common::hash::VirtualNode;
968 use risingwave_common::util::epoch::test_epoch;
969 use risingwave_hummock_sdk::key::UserKey;
970
971 use super::*;
972 use crate::assert_bytes_eq;
973 use crate::compaction_catalog_manager::{
974 CompactionCatalogAgent, DummyFilterKeyExtractor, MultiFilterKeyExtractor,
975 };
976 use crate::hummock::iterator::test_utils::mock_sstable_store;
977 use crate::hummock::sstable::xor_filter::BlockedXor16FilterBuilder;
978 use crate::hummock::test_utils::{
979 TEST_KEYS_COUNT, default_builder_opt_for_test, gen_test_sstable_impl, mock_sst_writer,
980 test_key_of, test_value_of,
981 };
982 use crate::hummock::{CachePolicy, Sstable, SstableWriterOptions, Xor8FilterBuilder};
983 use crate::monitor::StoreLocalStatistic;
984
985 #[tokio::test]
986 async fn test_empty() {
987 let opt = SstableBuilderOptions {
988 capacity: 0,
989 block_capacity: 4096,
990 restart_interval: 16,
991 bloom_false_positive: 0.001,
992 ..Default::default()
993 };
994
995 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
996 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
997 let b = SstableBuilder::for_test(
998 0,
999 mock_sst_writer(&opt),
1000 opt,
1001 table_id_to_vnode,
1002 table_id_to_watermark_serde,
1003 );
1004
1005 b.finish().await.unwrap();
1006 }
1007
1008 fn encode_full_key(vnode: VirtualNode, table_key_suffix: &[u8]) -> Vec<u8> {
1009 let mut table_key = vnode.to_be_bytes().to_vec();
1010 table_key.extend_from_slice(table_key_suffix);
1011 FullKey::for_test(TableId::default(), table_key, 0).encode()
1012 }
1013
1014 fn table_key_of(vnode: VirtualNode, suffix: &[u8]) -> Vec<u8> {
1015 let mut key = vnode.to_be_bytes().to_vec();
1016 key.extend_from_slice(suffix);
1017 key
1018 }
1019
1020 #[test]
1021 fn test_vnode_user_key_range_basic_collection() {
1022 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1025 let vnode_1 = VirtualNode::from_index(1);
1026 let vnode_2 = VirtualNode::from_index(2);
1027
1028 let k1 = encode_full_key(vnode_1, b"k1");
1029 let k2 = encode_full_key(vnode_1, b"k2");
1030 let k3 = encode_full_key(vnode_2, b"k3");
1031 let k4 = encode_full_key(vnode_2, b"k4");
1032
1033 collector.observe_key(vnode_1, &k1, &[]);
1034 collector.observe_key(vnode_1, &k2, &k1);
1035 collector.observe_key(vnode_2, &k3, &k2);
1036 collector.observe_key(vnode_2, &k4, &k3);
1037
1038 let info = collector.finish(&k4).unwrap();
1039 assert_eq!(info.vnode_user_key_ranges().len(), 2);
1040
1041 let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1043 assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1044 assert_eq!(
1045 range1_right.table_key.as_ref(),
1046 table_key_of(vnode_1, b"k2")
1047 );
1048
1049 let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1051 assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1052 assert_eq!(
1053 range2_right.table_key.as_ref(),
1054 table_key_of(vnode_2, b"k4")
1055 );
1056 }
1057
1058 #[test]
1059 fn test_vnode_user_key_range_capacity_limit() {
1060 let vnode_1 = VirtualNode::from_index(1);
1070 let vnode_2 = VirtualNode::from_index(2);
1071 let vnode_3 = VirtualNode::from_index(3);
1072 let vnode_4 = VirtualNode::from_index(4);
1073
1074 let k1 = encode_full_key(vnode_1, b"k1");
1075 let k2 = encode_full_key(vnode_1, b"k2");
1076 let k3 = encode_full_key(vnode_2, b"k3");
1077 let k4 = encode_full_key(vnode_2, b"k4");
1078 let k5 = encode_full_key(vnode_3, b"k5");
1079 let k6 = encode_full_key(vnode_3, b"k6");
1080 let k7 = encode_full_key(vnode_4, b"k7");
1081
1082 let range_size = mem::size_of::<VirtualNode>() + k1.len() + k2.len();
1083 let estimated_next_size = mem::size_of::<VirtualNode>() + k3.len();
1084 let limit = range_size + estimated_next_size; let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(limit)).unwrap();
1086
1087 collector.observe_key(vnode_1, &k1, &[]);
1088 collector.observe_key(vnode_1, &k2, &k1);
1089 collector.observe_key(vnode_2, &k3, &k2); collector.observe_key(vnode_2, &k4, &k3);
1091 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();
1096 assert_eq!(
1097 info.vnode_user_key_ranges().len(),
1098 2,
1099 "Should collect exactly 2 vnodes"
1100 );
1101
1102 let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1104 assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1105 assert_eq!(
1106 range1_right.table_key.as_ref(),
1107 table_key_of(vnode_1, b"k2")
1108 );
1109
1110 let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1111 assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1112 assert_eq!(
1113 range2_right.table_key.as_ref(),
1114 table_key_of(vnode_2, b"k4")
1115 );
1116
1117 assert!(info.get_vnode_user_key_range(vnode_3).is_none());
1119 assert!(info.get_vnode_user_key_range(vnode_4).is_none());
1120 }
1121
1122 #[test]
1123 fn test_vnode_user_key_range_sparse_distribution() {
1124 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1127 let vnode_5 = VirtualNode::from_index(5);
1128 let vnode_10 = VirtualNode::from_index(10);
1129 let vnode_100 = VirtualNode::from_index(100);
1130
1131 let keys = vec![
1132 (vnode_5, encode_full_key(vnode_5, b"key_005_001")),
1133 (vnode_5, encode_full_key(vnode_5, b"key_005_002")),
1134 (vnode_5, encode_full_key(vnode_5, b"key_005_999")),
1135 (vnode_10, encode_full_key(vnode_10, b"key_010_001")),
1136 (vnode_10, encode_full_key(vnode_10, b"key_010_100")),
1137 (vnode_100, encode_full_key(vnode_100, b"key_100_001")),
1138 (vnode_100, encode_full_key(vnode_100, b"key_100_999")),
1139 ];
1140
1141 let mut prev_key = Vec::new();
1142 for (vnode, key) in &keys {
1143 collector.observe_key(*vnode, key, &prev_key);
1144 prev_key = key.clone();
1145 }
1146
1147 let info = collector.finish(&prev_key).unwrap();
1148 assert_eq!(info.vnode_user_key_ranges().len(), 3);
1149
1150 let (range5_left, range5_right) = info.get_vnode_user_key_range(vnode_5).unwrap();
1152 assert_eq!(
1153 range5_left.table_key.as_ref(),
1154 table_key_of(vnode_5, b"key_005_001")
1155 );
1156 assert_eq!(
1157 range5_right.table_key.as_ref(),
1158 table_key_of(vnode_5, b"key_005_999")
1159 );
1160
1161 let (range10_left, range10_right) = info.get_vnode_user_key_range(vnode_10).unwrap();
1162 assert_eq!(
1163 range10_left.table_key.as_ref(),
1164 table_key_of(vnode_10, b"key_010_001")
1165 );
1166 assert_eq!(
1167 range10_right.table_key.as_ref(),
1168 table_key_of(vnode_10, b"key_010_100")
1169 );
1170
1171 let (range100_left, range100_right) = info.get_vnode_user_key_range(vnode_100).unwrap();
1172 assert_eq!(
1173 range100_left.table_key.as_ref(),
1174 table_key_of(vnode_100, b"key_100_001")
1175 );
1176 assert_eq!(
1177 range100_right.table_key.as_ref(),
1178 table_key_of(vnode_100, b"key_100_999")
1179 );
1180
1181 assert!(
1183 info.get_vnode_user_key_range(VirtualNode::from_index(0))
1184 .is_none()
1185 );
1186 assert!(
1187 info.get_vnode_user_key_range(VirtualNode::from_index(7))
1188 .is_none()
1189 );
1190 assert!(
1191 info.get_vnode_user_key_range(VirtualNode::from_index(50))
1192 .is_none()
1193 );
1194 }
1195
1196 #[test]
1197 fn test_vnode_user_key_range_edge_cases() {
1198 assert!(VnodeUserKeyRangeCollector::with_limit(None).is_none());
1200 assert!(VnodeUserKeyRangeCollector::with_limit(Some(0)).is_none());
1201
1202 let collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1204 assert!(
1205 collector
1206 .finish(&encode_full_key(VirtualNode::ZERO, b"any"))
1207 .is_none()
1208 );
1209
1210 let vnode = VirtualNode::from_index(5);
1212 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1213 let a = encode_full_key(vnode, b"a");
1214 let b = encode_full_key(vnode, b"b");
1215 let c = encode_full_key(vnode, b"c");
1216 collector.observe_key(vnode, &a, &[]);
1217 collector.observe_key(vnode, &b, &a);
1218 collector.observe_key(vnode, &c, &b);
1219 assert!(
1220 collector.finish(&c).is_none(),
1221 "Single-vnode SST should not emit hints"
1222 );
1223
1224 let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1)).unwrap();
1229 let vnode_1 = VirtualNode::from_index(1);
1230 let vnode_2 = VirtualNode::from_index(2);
1231 let key1 = encode_full_key(vnode_1, b"key1");
1232 let key2 = encode_full_key(vnode_1, b"key2");
1233 let key3 = encode_full_key(vnode_2, b"key3");
1234 collector.observe_key(vnode_1, &key1, &[]);
1235 collector.observe_key(vnode_1, &key2, &key1);
1236 collector.observe_key(vnode_2, &key3, &key2); assert!(
1238 collector.finish(&key3).is_none(),
1239 "Only 1 vnode collected, should return None"
1240 );
1241 }
1242
1243 #[tokio::test]
1244 async fn test_basic() {
1245 let opt = default_builder_opt_for_test();
1246
1247 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1248 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1249 let mut b = SstableBuilder::for_test(
1250 0,
1251 mock_sst_writer(&opt),
1252 opt,
1253 table_id_to_vnode,
1254 table_id_to_watermark_serde,
1255 );
1256
1257 for i in 0..TEST_KEYS_COUNT {
1258 b.add_for_test(
1259 test_key_of(i).to_ref(),
1260 HummockValue::put(&test_value_of(i)),
1261 )
1262 .await
1263 .unwrap();
1264 }
1265
1266 let output = b.finish().await.unwrap();
1267 let info = output.sst_info.sst_info;
1268
1269 assert_bytes_eq!(test_key_of(0).encode(), info.key_range.left);
1270 assert_bytes_eq!(
1271 test_key_of(TEST_KEYS_COUNT - 1).encode(),
1272 info.key_range.right
1273 );
1274 let (data, meta) = output.writer_output;
1275 assert_eq!(info.file_size, meta.estimated_size as u64);
1276 let offset = info.meta_offset as usize;
1277 let meta2 = SstableMeta::decode(&data[offset..]).unwrap();
1278 assert_eq!(meta2, meta);
1279 }
1280
1281 async fn test_with_xor_filter_builder<F: FilterBuilder>(
1282 bloom_false_positive: f64,
1283 expected_filter_type: PbSstableFilterType,
1284 expected_filter_layout: PbSstableFilterLayout,
1285 ) {
1286 let key_count = 1000;
1287
1288 let opts = SstableBuilderOptions {
1289 capacity: 0,
1290 block_capacity: 4096,
1291 restart_interval: 16,
1292 bloom_false_positive,
1293 ..Default::default()
1294 };
1295
1296 let sstable_store = mock_sstable_store().await;
1298 let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1299 let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1300 let sst_info = gen_test_sstable_impl::<Vec<u8>, F>(
1301 opts,
1302 0,
1303 (0..TEST_KEYS_COUNT).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
1304 sstable_store.clone(),
1305 CachePolicy::NotFill,
1306 table_id_to_vnode,
1307 table_id_to_watermark_serde,
1308 )
1309 .await;
1310 assert_eq!(sst_info.filter_type, expected_filter_type);
1311 assert_eq!(sst_info.filter_layout, expected_filter_layout);
1312 let table = sstable_store
1313 .sstable(&sst_info, &mut StoreLocalStatistic::default())
1314 .await
1315 .unwrap();
1316
1317 assert!(table.has_filter());
1318 for i in 0..key_count {
1319 let full_key = test_key_of(i);
1320 let hash = Sstable::hash_for_filter(full_key.user_key.encode().as_slice(), 0);
1321 let key_ref = full_key.user_key.as_ref();
1322 assert!(
1323 table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash),
1324 "failed at {}",
1325 i
1326 );
1327 }
1328 }
1329
1330 #[tokio::test]
1331 async fn test_xor_filter_builder_output() {
1332 test_with_xor_filter_builder::<Xor16FilterBuilder>(
1333 0.0,
1334 PbSstableFilterType::SstableFilterXor16,
1335 PbSstableFilterLayout::Plain,
1336 )
1337 .await;
1338 test_with_xor_filter_builder::<Xor16FilterBuilder>(
1339 0.01,
1340 PbSstableFilterType::SstableFilterXor16,
1341 PbSstableFilterLayout::Plain,
1342 )
1343 .await;
1344 test_with_xor_filter_builder::<Xor8FilterBuilder>(
1345 0.01,
1346 PbSstableFilterType::SstableFilterXor8,
1347 PbSstableFilterLayout::Plain,
1348 )
1349 .await;
1350 test_with_xor_filter_builder::<BlockedXor16FilterBuilder>(
1351 0.01,
1352 PbSstableFilterType::SstableFilterXor16,
1353 PbSstableFilterLayout::Blocked,
1354 )
1355 .await;
1356 }
1357
1358 #[tokio::test]
1359 async fn test_no_xor_filter_block() {
1360 let opts = SstableBuilderOptions::default();
1361 let sstable_store = mock_sstable_store().await;
1363 let writer_opts = SstableWriterOptions::default();
1364 let object_id = 1;
1365 let writer = sstable_store
1366 .clone()
1367 .create_sst_writer(object_id, writer_opts);
1368 let mut filter = MultiFilterKeyExtractor::default();
1369 filter.register(
1370 1.into(),
1371 FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1372 );
1373 filter.register(
1374 2.into(),
1375 FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
1376 );
1377 filter.register(
1378 3.into(),
1379 FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1380 );
1381
1382 let table_id_to_vnode = HashMap::from_iter(vec![
1383 (1.into(), VirtualNode::COUNT_FOR_TEST),
1384 (2.into(), VirtualNode::COUNT_FOR_TEST),
1385 (3.into(), VirtualNode::COUNT_FOR_TEST),
1386 ]);
1387 let table_id_to_watermark_serde =
1388 HashMap::from_iter(vec![(1.into(), None), (2.into(), None), (3.into(), None)]);
1389
1390 let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
1391 FilterKeyExtractorImpl::Multi(filter),
1392 table_id_to_vnode,
1393 table_id_to_watermark_serde,
1394 HashMap::default(),
1395 ));
1396
1397 let mut builder = SstableBuilder::new(
1398 object_id,
1399 writer,
1400 BlockedXor16FilterBuilder::create(opts.filter_builder_options()),
1401 opts,
1402 compaction_catalog_agent_ref,
1403 None,
1404 );
1405
1406 let key_count: usize = 10000;
1407 for table_id in 1..4 {
1408 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1409 for idx in 0..key_count {
1410 table_key.resize(VirtualNode::SIZE, 0);
1411 table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1412 let k = UserKey::for_test(TableId::new(table_id), table_key.as_ref());
1413 let v = test_value_of(idx);
1414 builder
1415 .add(
1416 FullKey::from_user_key(k, test_epoch(1)),
1417 HummockValue::put(v.as_ref()),
1418 )
1419 .await
1420 .unwrap();
1421 }
1422 }
1423 let ret = builder.finish().await.unwrap();
1424 let sst_info = ret.sst_info.sst_info.clone();
1425 ret.writer_output.await.unwrap().unwrap();
1426 let table = sstable_store
1427 .sstable(&sst_info, &mut StoreLocalStatistic::default())
1428 .await
1429 .unwrap();
1430 let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1431 for idx in 0..key_count {
1432 table_key.resize(VirtualNode::SIZE, 0);
1433 table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1434 let k = UserKey::for_test(TableId::new(2), table_key.as_slice());
1435 let hash = Sstable::hash_for_filter(&k.encode(), 2);
1436 let key_ref = k.as_ref();
1437 assert!(
1438 table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash)
1439 );
1440 }
1441 }
1442}