1use std::collections::HashMap;
20use std::time::Duration;
21
22use common_base::readable_size::ReadableSize;
23use common_telemetry::info;
24use common_time::TimeToLive;
25use common_wal::options::{WAL_OPTIONS_KEY, WalOptions};
26use serde::de::Error as _;
27use serde::{Deserialize, Deserializer, Serialize};
28use serde_json::Value;
29use serde_with::{DisplayFromStr, NoneAsEmptyString, serde_as, with_prefix};
30use snafu::{ResultExt, ensure};
31use store_api::codec::PrimaryKeyEncoding;
32use store_api::metric_engine_consts::{
33 MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING, PRIMARY_KEY_ENCODING,
34};
35use store_api::mito_engine_options::{COMPACTION_OVERRIDE, MAX_ROW_GROUP_ROW_COUNT_LIMIT};
36use store_api::storage::{ColumnId, RegionId};
37use strum::EnumString;
38
39use crate::error::{InvalidRegionOptionsSnafu, JsonOptionsSnafu, Result};
40use crate::memtable::bulk::BulkMemtableConfig;
41use crate::sst::FormatType;
42use crate::sst::parquet::DEFAULT_ROW_GROUP_SIZE;
43
44const DEFAULT_INDEX_SEGMENT_ROW_COUNT: usize = 1024;
45const COMPACTION_TWCS_PREFIX: &str = "compaction.twcs.";
46const MEMTABLE_PARTITION_TREE_PREFIX: &str = "memtable.partition_tree.";
47const MEMTABLE_BULK_PREFIX: &str = "memtable.bulk.";
48
49const LEGACY_PARTITION_TREE_MEMTABLE_TYPE: &str = "partition_tree";
53
54pub(crate) fn parse_wal_options(
55 options_map: &HashMap<String, String>,
56) -> std::result::Result<WalOptions, serde_json::Error> {
57 options_map
58 .get(WAL_OPTIONS_KEY)
59 .map_or(Ok(WalOptions::default()), |encoded_wal_options| {
60 serde_json::from_str(encoded_wal_options)
61 })
62}
63
64#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, EnumString)]
66#[serde(rename_all = "snake_case")]
67#[strum(serialize_all = "snake_case")]
68pub enum MergeMode {
69 #[default]
71 LastRow,
72 LastNonNull,
74}
75
76#[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
82#[serde(default)]
83pub struct RegionOptions {
84 pub ttl: Option<TimeToLive>,
86 #[serde(with = "humantime_serde")]
89 pub auto_flush_interval: Option<Duration>,
90 pub compaction: CompactionOptions,
92 pub compaction_override: bool,
93 pub storage: Option<String>,
95 pub append_mode: bool,
97 pub skip_wal: bool,
99 pub wal_options: WalOptions,
101 pub index_options: IndexOptions,
103 pub memtable: Option<MemtableOptions>,
105 pub merge_mode: Option<MergeMode>,
108 pub sst_format: Option<FormatType>,
110 #[serde(skip_serializing_if = "Option::is_none")]
112 pub max_row_group_row_count: Option<usize>,
113 #[serde(skip_serializing_if = "Option::is_none")]
115 pub primary_key_encoding: Option<PrimaryKeyEncoding>,
116 #[serde(skip_serializing_if = "Option::is_none")]
120 pub write_buffer_size: Option<ReadableSize>,
121}
122
123impl RegionOptions {
124 pub fn validate(&self) -> Result<()> {
126 if self.append_mode {
127 ensure!(
128 self.merge_mode
129 .is_none_or(|mode| mode == MergeMode::LastRow),
130 InvalidRegionOptionsSnafu {
131 reason: "only last_row merge_mode is allowed when append_mode is enabled",
132 }
133 );
134 }
135 if let Some(auto_flush_interval) = self.auto_flush_interval {
136 ensure!(
137 auto_flush_interval > Duration::ZERO,
138 InvalidRegionOptionsSnafu {
139 reason: "auto_flush_interval must be greater than 0",
140 }
141 );
142 }
143 if let Some(row_count) = self.max_row_group_row_count {
144 ensure!(
145 row_count > 0 && row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT,
146 InvalidRegionOptionsSnafu {
147 reason: format!(
148 "max_row_group_row_count must be in (0, {MAX_ROW_GROUP_ROW_COUNT_LIMIT}], got {row_count}",
149 ),
150 }
151 );
152 }
153 Ok(())
154 }
155
156 pub fn row_group_size(&self) -> usize {
158 self.max_row_group_row_count
159 .unwrap_or(DEFAULT_ROW_GROUP_SIZE)
160 }
161
162 pub fn need_dedup(&self) -> bool {
164 !self.append_mode
165 }
166
167 pub fn merge_mode(&self) -> MergeMode {
169 self.merge_mode.unwrap_or_default()
170 }
171
172 pub fn auto_flush_interval_or(&self, default: Duration) -> Duration {
174 self.auto_flush_interval.unwrap_or(default)
175 }
176
177 pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
179 self.primary_key_encoding.unwrap_or_default()
180 }
181}
182
183impl RegionOptions {
184 pub fn try_from_options(
186 region_id: RegionId,
187 options_map: &HashMap<String, String>,
188 ) -> Result<Self> {
189 let value = options_map_to_value(options_map);
190 let json = serde_json::to_string(&value).context(JsonOptionsSnafu)?;
191
192 let options: RegionOptionsWithoutEnum =
196 serde_json::from_str(&json).context(JsonOptionsSnafu)?;
197 let has_compaction_type =
198 validate_enum_options(options_map, "compaction.type", &[COMPACTION_TWCS_PREFIX])?;
199 let compaction = if has_compaction_type {
200 serde_json::from_str(&json).context(JsonOptionsSnafu)?
201 } else {
202 CompactionOptions::default()
203 };
204
205 let wal_options = parse_wal_options(options_map).context(JsonOptionsSnafu)?;
206
207 let index_options: IndexOptions = serde_json::from_str(&json).context(JsonOptionsSnafu)?;
208 let is_legacy_partition_tree = options_map
209 .get("memtable.type")
210 .map(|s| s.eq_ignore_ascii_case(LEGACY_PARTITION_TREE_MEMTABLE_TYPE))
211 .unwrap_or(false);
212 let memtable = if validate_enum_options(
213 options_map,
214 "memtable.type",
215 &[MEMTABLE_PARTITION_TREE_PREFIX, MEMTABLE_BULK_PREFIX],
216 )? {
217 if is_legacy_partition_tree {
218 None
222 } else {
223 Some(serde_json::from_str(&json).context(JsonOptionsSnafu)?)
224 }
225 } else {
226 None
227 };
228
229 let mut sst_format = options.sst_format;
232 if is_legacy_partition_tree {
233 info!(
234 "Region {} specified the removed partition_tree memtable; \
235 overriding memtable to the default and SST format to flat",
236 region_id
237 );
238 sst_format = Some(FormatType::Flat);
239 }
240
241 if matches!(memtable, Some(MemtableOptions::Bulk(_))) {
244 if let Some(format) = sst_format
245 && format != FormatType::Flat
246 {
247 info!(
248 "Region {} uses bulk memtable; overriding sst_format from {:?} to flat",
249 region_id, format
250 );
251 }
252 sst_format = Some(FormatType::Flat);
253 }
254
255 let compaction_override_flag = options_map
256 .get(COMPACTION_OVERRIDE)
257 .map(|v| matches!(v.to_lowercase().as_str(), "true" | "1"))
258 .unwrap_or(false);
259 let compaction_override = has_compaction_type || compaction_override_flag;
260 let primary_key_encoding = options_map
261 .get(PRIMARY_KEY_ENCODING)
262 .or_else(|| options_map.get(MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING))
263 .map(|v| match v.to_lowercase().as_str() {
264 "dense" => Ok(PrimaryKeyEncoding::Dense),
265 "sparse" => Ok(PrimaryKeyEncoding::Sparse),
266 _ => Err(InvalidRegionOptionsSnafu {
267 reason: format!("Invalid primary key encoding: {v}"),
268 }
269 .build()),
270 })
271 .transpose()?;
272
273 let opts = RegionOptions {
274 ttl: options.ttl,
275 auto_flush_interval: options.auto_flush_interval,
276 compaction,
277 compaction_override,
278 storage: options.storage,
279 append_mode: options.append_mode,
280 skip_wal: options.skip_wal,
281 wal_options,
282 index_options,
283 memtable,
284 merge_mode: options.merge_mode,
285 sst_format,
286 max_row_group_row_count: options.max_row_group_row_count,
287 primary_key_encoding,
288 write_buffer_size: options.write_buffer_size,
289 };
290 opts.validate()?;
291
292 Ok(opts)
293 }
294}
295
296#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
298#[serde(tag = "compaction.type")]
299#[serde(rename_all = "snake_case")]
300pub enum CompactionOptions {
301 #[serde(with = "prefix_twcs")]
303 Twcs(TwcsOptions),
304}
305
306impl CompactionOptions {
307 pub(crate) fn time_window(&self) -> Option<Duration> {
308 match self {
309 CompactionOptions::Twcs(opts) => opts.time_window,
310 }
311 }
312
313 pub(crate) fn remote_compaction(&self) -> bool {
314 match self {
315 CompactionOptions::Twcs(opts) => opts.remote_compaction,
316 }
317 }
318
319 pub(crate) fn fallback_to_local(&self) -> bool {
320 match self {
321 CompactionOptions::Twcs(opts) => opts.fallback_to_local,
322 }
323 }
324}
325
326impl Default for CompactionOptions {
327 fn default() -> Self {
328 Self::Twcs(TwcsOptions::default())
329 }
330}
331
332#[serde_as]
334#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
335#[serde(default)]
336pub struct TwcsOptions {
337 #[serde_as(as = "DisplayFromStr")]
339 pub trigger_file_num: usize,
340 #[serde(with = "humantime_serde")]
342 pub time_window: Option<Duration>,
343 pub max_output_file_size: Option<ReadableSize>,
345 #[serde_as(as = "DisplayFromStr")]
347 pub remote_compaction: bool,
348 #[serde_as(as = "DisplayFromStr")]
350 pub fallback_to_local: bool,
351}
352
353with_prefix!(prefix_twcs "compaction.twcs.");
354
355impl TwcsOptions {
356 pub fn time_window_seconds(&self) -> Option<i64> {
358 self.time_window.and_then(|window| {
359 let window_secs = window.as_secs();
360 if window_secs == 0 {
361 None
362 } else {
363 window_secs.try_into().ok()
364 }
365 })
366 }
367}
368
369impl Default for TwcsOptions {
370 fn default() -> Self {
371 Self {
372 trigger_file_num: 4,
373 time_window: None,
374 max_output_file_size: Some(ReadableSize::mb(512)),
375 remote_compaction: false,
376 fallback_to_local: true,
377 }
378 }
379}
380
381#[serde_as]
384#[derive(Debug, Deserialize)]
385#[serde(default)]
386struct RegionOptionsWithoutEnum {
387 write_buffer_size: Option<ReadableSize>,
388 ttl: Option<TimeToLive>,
390 #[serde(with = "humantime_serde")]
391 auto_flush_interval: Option<Duration>,
392 storage: Option<String>,
393 #[serde_as(as = "DisplayFromStr")]
394 append_mode: bool,
395 #[serde_as(as = "DisplayFromStr")]
396 skip_wal: bool,
397 #[serde_as(as = "NoneAsEmptyString")]
398 merge_mode: Option<MergeMode>,
399 #[serde_as(as = "NoneAsEmptyString")]
400 sst_format: Option<FormatType>,
401 #[serde_as(as = "NoneAsEmptyString")]
402 max_row_group_row_count: Option<usize>,
403}
404
405impl Default for RegionOptionsWithoutEnum {
406 fn default() -> Self {
407 let options = RegionOptions::default();
408 RegionOptionsWithoutEnum {
409 write_buffer_size: options.write_buffer_size,
410 ttl: options.ttl,
411 auto_flush_interval: options.auto_flush_interval,
412 storage: options.storage,
413 append_mode: options.append_mode,
414 skip_wal: options.skip_wal,
415 merge_mode: options.merge_mode,
416 sst_format: options.sst_format,
417 max_row_group_row_count: options.max_row_group_row_count,
418 }
419 }
420}
421
422with_prefix!(prefix_inverted_index "index.inverted_index.");
423
424#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
426#[serde(default)]
427pub struct IndexOptions {
428 #[serde(flatten, with = "prefix_inverted_index")]
430 pub inverted_index: InvertedIndexOptions,
431}
432
433#[serde_as]
435#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
436#[serde(default)]
437pub struct InvertedIndexOptions {
438 #[serde(deserialize_with = "deserialize_ignore_column_ids")]
441 #[serde(serialize_with = "serialize_ignore_column_ids")]
442 pub ignore_column_ids: Vec<ColumnId>,
443
444 #[serde_as(as = "DisplayFromStr")]
446 pub segment_row_count: usize,
447}
448
449impl Default for InvertedIndexOptions {
450 fn default() -> Self {
451 Self {
452 ignore_column_ids: Vec::new(),
453 segment_row_count: DEFAULT_INDEX_SEGMENT_ROW_COUNT,
454 }
455 }
456}
457
458#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
460#[serde(tag = "memtable.type", rename_all = "snake_case")]
461pub enum MemtableOptions {
462 TimeSeries,
463 #[serde(with = "prefix_bulk")]
464 Bulk(BulkMemtableConfig),
465}
466
467with_prefix!(prefix_bulk "memtable.bulk.");
468
469fn deserialize_ignore_column_ids<'de, D>(deserializer: D) -> Result<Vec<ColumnId>, D::Error>
470where
471 D: Deserializer<'de>,
472{
473 let s: String = Deserialize::deserialize(deserializer)?;
474 let mut column_ids = Vec::new();
475 if s.is_empty() {
476 return Ok(column_ids);
477 }
478 for item in s.split(',') {
479 let column_id = item.parse().map_err(D::Error::custom)?;
480 column_ids.push(column_id);
481 }
482 Ok(column_ids)
483}
484
485fn serialize_ignore_column_ids<S>(column_ids: &[ColumnId], serializer: S) -> Result<S::Ok, S::Error>
486where
487 S: serde::Serializer,
488{
489 let s = column_ids
490 .iter()
491 .map(|id| id.to_string())
492 .collect::<Vec<_>>()
493 .join(",");
494 serializer.serialize_str(&s)
495}
496
497fn options_map_to_value(options: &HashMap<String, String>) -> Value {
501 let map = options
502 .iter()
503 .map(|(key, value)| {
504 if value.eq_ignore_ascii_case("null") {
506 (key.clone(), Value::Null)
507 } else {
508 (key.clone(), Value::from(value.clone()))
509 }
510 })
511 .collect();
512 Value::Object(map)
513}
514
515fn validate_enum_options(
523 options_map: &HashMap<String, String>,
524 enum_tag_key: &str,
525 enum_option_prefixes: &[&str],
526) -> Result<bool> {
527 let mut has_enum_options = false;
528 let mut has_tag = false;
529 for key in options_map.keys() {
530 if key == enum_tag_key {
531 has_tag = true;
532 } else if !has_enum_options
533 && enum_option_prefixes
534 .iter()
535 .any(|prefix| key.starts_with(prefix))
536 {
537 has_enum_options = true;
538 }
539
540 if has_tag && has_enum_options {
541 break;
542 }
543 }
544
545 ensure!(
547 has_tag || !has_enum_options,
548 InvalidRegionOptionsSnafu {
549 reason: format!("missing key {} in options", enum_tag_key),
550 }
551 );
552
553 Ok(has_tag)
554}
555
556#[cfg(test)]
557mod tests {
558 use common_error::ext::ErrorExt;
559 use common_error::status_code::StatusCode;
560 use common_wal::options::KafkaWalOptions;
561 use store_api::mito_engine_options::{SKIP_WAL_KEY, WRITE_BUFFER_SIZE_KEY};
562
563 use super::*;
564
565 fn make_map(options: &[(&str, &str)]) -> HashMap<String, String> {
566 options
567 .iter()
568 .map(|(k, v)| (k.to_string(), v.to_string()))
569 .collect()
570 }
571
572 #[test]
573 fn test_empty_region_options() {
574 let map = make_map(&[]);
575 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
576 assert_eq!(RegionOptions::default(), options);
577 }
578
579 #[test]
580 fn test_with_ttl() {
581 let map = make_map(&[("ttl", "7d")]);
582 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
583 let expect = RegionOptions {
584 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
585 ..Default::default()
586 };
587 assert_eq!(expect, options);
588 }
589
590 #[test]
591 fn test_with_skip_wal() {
592 let map = make_map(&[(SKIP_WAL_KEY, "true")]);
593 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
594 assert!(options.skip_wal);
595 }
596
597 #[test]
598 fn test_with_auto_flush_interval() {
599 let map = make_map(&[("auto_flush_interval", "5m")]);
600 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
601 let expect = RegionOptions {
602 auto_flush_interval: Some(Duration::from_secs(5 * 60)),
603 ..Default::default()
604 };
605 assert_eq!(expect, options);
606 }
607
608 #[test]
609 fn test_with_zero_auto_flush_interval() {
610 let map = make_map(&[("auto_flush_interval", "0s")]);
611 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
612 assert!(
613 err.to_string().contains("auto_flush_interval"),
614 "unexpected error: {err}"
615 );
616 }
617
618 #[test]
619 fn test_with_write_buffer_size() {
620 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "128MiB")]);
621 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
622 let expect = RegionOptions {
623 write_buffer_size: Some(ReadableSize::mb(128)),
624 ..Default::default()
625 };
626 assert_eq!(expect, options);
627 }
628
629 #[test]
630 fn test_with_zero_write_buffer_size() {
631 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "0")]);
632 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
633 let expect = RegionOptions {
634 write_buffer_size: Some(ReadableSize::mb(0)),
635 ..Default::default()
636 };
637 assert_eq!(expect, options);
638 }
639
640 #[test]
641 fn test_with_storage() {
642 let map = make_map(&[("storage", "S3")]);
643 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
644 let expect = RegionOptions {
645 storage: Some("S3".to_string()),
646 ..Default::default()
647 };
648 assert_eq!(expect, options);
649 }
650
651 #[test]
652 fn test_without_compaction_type() {
653 let map = make_map(&[
654 ("compaction.twcs.trigger_file_num", "8"),
655 ("compaction.twcs.time_window", "2h"),
656 ]);
657 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
658 assert_eq!(StatusCode::InvalidArguments, err.status_code());
659 }
660
661 #[test]
662 fn test_with_compaction_type() {
663 let map = make_map(&[
664 ("compaction.twcs.trigger_file_num", "8"),
665 ("compaction.twcs.time_window", "2h"),
666 ("compaction.type", "twcs"),
667 ]);
668 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
669 let expect = RegionOptions {
670 compaction: CompactionOptions::Twcs(TwcsOptions {
671 trigger_file_num: 8,
672 time_window: Some(Duration::from_secs(3600 * 2)),
673 ..Default::default()
674 }),
675 compaction_override: true,
676 ..Default::default()
677 };
678 assert_eq!(expect, options);
679 }
680
681 #[test]
682 fn test_with_compaction_override_true_without_compaction_type() {
683 let map = make_map(&[(COMPACTION_OVERRIDE, "true")]);
684 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
685 let expect = RegionOptions {
686 compaction_override: true,
687 ..Default::default()
688 };
689 assert_eq!(expect, options);
690 }
691
692 #[test]
693 fn test_with_compaction_override_false_without_compaction_type() {
694 let map = make_map(&[(COMPACTION_OVERRIDE, "false")]);
695 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
696 assert_eq!(RegionOptions::default(), options);
697 }
698
699 #[test]
700 fn test_compaction_twcs_options_still_require_compaction_type_with_override() {
701 let map = make_map(&[
702 (COMPACTION_OVERRIDE, "true"),
703 ("compaction.twcs.time_window", "2h"),
704 ]);
705 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
706 assert_eq!(StatusCode::InvalidArguments, err.status_code());
707 }
708
709 fn test_with_wal_options(wal_options: &WalOptions) -> bool {
710 let encoded_wal_options = serde_json::to_string(&wal_options).unwrap();
711 let map = make_map(&[(WAL_OPTIONS_KEY, &encoded_wal_options)]);
712 let got = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
713 let expect = RegionOptions {
714 wal_options: wal_options.clone(),
715 ..Default::default()
716 };
717 expect == got
718 }
719
720 #[test]
721 fn test_with_index() {
722 let map = make_map(&[
723 ("index.inverted_index.ignore_column_ids", "1,2,3"),
724 ("index.inverted_index.segment_row_count", "512"),
725 ]);
726 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
727 let expect = RegionOptions {
728 index_options: IndexOptions {
729 inverted_index: InvertedIndexOptions {
730 ignore_column_ids: vec![1, 2, 3],
731 segment_row_count: 512,
732 },
733 },
734 ..Default::default()
735 };
736 assert_eq!(expect, options);
737 }
738
739 #[test]
741 fn test_with_any_wal_options() {
742 let all_wal_options = [
743 WalOptions::RaftEngine,
744 WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
745 ];
746 all_wal_options.iter().all(test_with_wal_options);
747 }
748
749 #[test]
750 fn test_with_memtable() {
751 let map = make_map(&[("memtable.type", "time_series")]);
752 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
753 let expect = RegionOptions {
754 memtable: Some(MemtableOptions::TimeSeries),
755 ..Default::default()
756 };
757 assert_eq!(expect, options);
758
759 let map = make_map(&[("memtable.type", "bulk")]);
760 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
761 let expect = RegionOptions {
762 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
763 sst_format: Some(FormatType::Flat),
764 ..Default::default()
765 };
766 assert_eq!(expect, options);
767
768 let map = make_map(&[
769 ("memtable.type", "bulk"),
770 ("memtable.bulk.merge_threshold", "7"),
771 ("memtable.bulk.encode_row_threshold", "11"),
772 ("memtable.bulk.encode_bytes_threshold", "13"),
773 ("memtable.bulk.max_merge_groups", "17"),
774 ]);
775 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
776 let expect = RegionOptions {
777 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
778 merge_threshold: 7,
779 encode_row_threshold: 11,
780 encode_bytes_threshold: 13,
781 max_merge_groups: 17,
782 })),
783 sst_format: Some(FormatType::Flat),
784 ..Default::default()
785 };
786 assert_eq!(expect, options);
787
788 let map = make_map(&[("memtable.type", "partition_tree")]);
791 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
792 let expect = RegionOptions {
793 memtable: None,
794 sst_format: Some(FormatType::Flat),
795 ..Default::default()
796 };
797 assert_eq!(expect, options);
798
799 let map = make_map(&[
801 ("memtable.type", "partition_tree"),
802 ("memtable.partition_tree.index_max_keys_per_shard", "2048"),
803 ("memtable.partition_tree.fork_dictionary_bytes", "128M"),
804 ]);
805 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
806 let expect = RegionOptions {
807 memtable: None,
808 sst_format: Some(FormatType::Flat),
809 ..Default::default()
810 };
811 assert_eq!(expect, options);
812 }
813
814 #[test]
815 fn test_primary_key_encoding() {
816 let map = make_map(&[("primary_key_encoding", "sparse")]);
818 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
819 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
820 assert_eq!(
821 options.primary_key_encoding,
822 Some(PrimaryKeyEncoding::Sparse)
823 );
824
825 let map = make_map(&[
827 ("memtable.type", "partition_tree"),
828 ("memtable.partition_tree.primary_key_encoding", "sparse"),
829 ]);
830 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
831 assert_eq!(options.memtable, None);
832 assert_eq!(options.sst_format, Some(FormatType::Flat));
833 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
834
835 let map = make_map(&[("primary_key_encoding", "bogus")]);
837 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
838 assert_eq!(StatusCode::InvalidArguments, err.status_code());
839 }
840
841 #[test]
842 fn test_legacy_partition_tree_overrides_sst_format() {
843 let map = make_map(&[
846 ("memtable.type", "partition_tree"),
847 ("sst_format", "primary_key"),
848 ]);
849 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
850 assert_eq!(options.memtable, None);
851 assert_eq!(options.sst_format, Some(FormatType::Flat));
852 }
853
854 #[test]
855 fn test_bulk_memtable_overrides_sst_format() {
856 let map = make_map(&[("memtable.type", "bulk"), ("sst_format", "primary_key")]);
860 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
861 assert_eq!(
862 options.memtable,
863 Some(MemtableOptions::Bulk(BulkMemtableConfig::default()))
864 );
865 assert_eq!(options.sst_format, Some(FormatType::Flat));
866 }
867
868 #[test]
869 fn test_unknown_memtable_type() {
870 let map = make_map(&[("memtable.type", "no_such_memtable")]);
871 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
872 assert_eq!(StatusCode::InvalidArguments, err.status_code());
873 }
874
875 #[test]
876 fn test_without_memtable_type() {
877 let map = make_map(&[("memtable.partition_tree.index_max_keys_per_shard", "2048")]);
878 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
879 assert_eq!(StatusCode::InvalidArguments, err.status_code());
880
881 let map = make_map(&[("memtable.bulk.merge_threshold", "7")]);
882 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
883 assert_eq!(StatusCode::InvalidArguments, err.status_code());
884 }
885
886 #[test]
887 fn test_with_merge_mode() {
888 let map = make_map(&[("merge_mode", "last_row")]);
889 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
890 assert_eq!(MergeMode::LastRow, options.merge_mode());
891
892 let map = make_map(&[("merge_mode", "last_non_null")]);
893 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
894 assert_eq!(MergeMode::LastNonNull, options.merge_mode());
895
896 let map = make_map(&[("merge_mode", "unknown")]);
897 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
898 assert_eq!(StatusCode::InvalidArguments, err.status_code());
899 }
900
901 #[test]
902 fn test_append_mode_allows_last_row_merge_mode() {
903 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_row")]);
904 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
905 assert!(options.append_mode);
906 assert_eq!(MergeMode::LastRow, options.merge_mode());
907
908 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_non_null")]);
909 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
910 assert_eq!(StatusCode::InvalidArguments, err.status_code());
911 }
912
913 #[test]
914 fn test_with_all() {
915 let wal_options = WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string()));
916 let map = make_map(&[
917 ("ttl", "7d"),
918 ("compaction.twcs.trigger_file_num", "8"),
919 ("compaction.twcs.max_output_file_size", "1GB"),
920 ("compaction.twcs.time_window", "2h"),
921 ("compaction.type", "twcs"),
922 ("compaction.twcs.remote_compaction", "false"),
923 ("compaction.twcs.fallback_to_local", "true"),
924 ("storage", "S3"),
925 ("append_mode", "false"),
926 ("index.inverted_index.ignore_column_ids", "1,2,3"),
927 ("index.inverted_index.segment_row_count", "512"),
928 (
929 WAL_OPTIONS_KEY,
930 &serde_json::to_string(&wal_options).unwrap(),
931 ),
932 ("memtable.type", "bulk"),
933 ("memtable.bulk.merge_threshold", "7"),
934 ("memtable.bulk.encode_row_threshold", "11"),
935 ("memtable.bulk.encode_bytes_threshold", "13"),
936 ("memtable.bulk.max_merge_groups", "17"),
937 ("merge_mode", "last_non_null"),
938 ]);
939 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
940 let expect = RegionOptions {
941 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
942 auto_flush_interval: None,
943 compaction: CompactionOptions::Twcs(TwcsOptions {
944 trigger_file_num: 8,
945 time_window: Some(Duration::from_secs(3600 * 2)),
946 max_output_file_size: Some(ReadableSize::gb(1)),
947 remote_compaction: false,
948 fallback_to_local: true,
949 }),
950 compaction_override: true,
951 storage: Some("S3".to_string()),
952 append_mode: false,
953 skip_wal: false,
954 wal_options,
955 index_options: IndexOptions {
956 inverted_index: InvertedIndexOptions {
957 ignore_column_ids: vec![1, 2, 3],
958 segment_row_count: 512,
959 },
960 },
961 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
962 merge_threshold: 7,
963 encode_row_threshold: 11,
964 encode_bytes_threshold: 13,
965 max_merge_groups: 17,
966 })),
967 merge_mode: Some(MergeMode::LastNonNull),
968 sst_format: Some(FormatType::Flat),
969 max_row_group_row_count: None,
970 primary_key_encoding: None,
971 write_buffer_size: None,
972 };
973 assert_eq!(expect, options);
974 }
975
976 #[test]
977 fn test_region_options_serde() {
978 let options = RegionOptions {
979 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
980 auto_flush_interval: None,
981 compaction: CompactionOptions::Twcs(TwcsOptions {
982 trigger_file_num: 8,
983 time_window: Some(Duration::from_secs(3600 * 2)),
984 max_output_file_size: None,
985 remote_compaction: false,
986 fallback_to_local: true,
987 }),
988 compaction_override: false,
989 storage: Some("S3".to_string()),
990 append_mode: false,
991 skip_wal: false,
992 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
993 index_options: IndexOptions {
994 inverted_index: InvertedIndexOptions {
995 ignore_column_ids: vec![1, 2, 3],
996 segment_row_count: 512,
997 },
998 },
999 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1000 merge_mode: Some(MergeMode::LastNonNull),
1001 sst_format: None,
1002 max_row_group_row_count: None,
1003 primary_key_encoding: None,
1004 write_buffer_size: Some(ReadableSize::mb(128)),
1005 };
1006 let region_options_json_str = serde_json::to_string(&options).unwrap();
1007 let got: RegionOptions = serde_json::from_str(®ion_options_json_str).unwrap();
1008 assert_eq!(options, got);
1009
1010 let old_region_options_json_str = r#"{"ttl":null}"#;
1011 let got: RegionOptions = serde_json::from_str(old_region_options_json_str).unwrap();
1012 assert_eq!(None, got.write_buffer_size);
1013
1014 let default_json = serde_json::to_value(RegionOptions::default()).unwrap();
1015 assert!(default_json.get(WRITE_BUFFER_SIZE_KEY).is_none());
1016 }
1017
1018 #[test]
1019 fn test_region_options_str_serde() {
1020 let region_options_json_str = r#"{
1022 "ttl": "7days",
1023 "compaction": {
1024 "compaction.type": "twcs",
1025 "compaction.twcs.trigger_file_num": "8",
1026 "compaction.twcs.max_output_file_size": "7MB",
1027 "compaction.twcs.time_window": "2h"
1028 },
1029 "storage": "S3",
1030 "append_mode": false,
1031 "wal_options": {
1032 "wal.provider": "kafka",
1033 "wal.kafka.topic": "test_topic"
1034 },
1035 "index_options": {
1036 "index.inverted_index.ignore_column_ids": "",
1037 "index.inverted_index.segment_row_count": "512"
1038 },
1039 "memtable": {
1040 "memtable.type": "bulk"
1041 },
1042 "merge_mode": "last_non_null"
1043}"#;
1044 let got: RegionOptions = serde_json::from_str(region_options_json_str).unwrap();
1045 let options = RegionOptions {
1046 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1047 auto_flush_interval: None,
1048 compaction: CompactionOptions::Twcs(TwcsOptions {
1049 trigger_file_num: 8,
1050 time_window: Some(Duration::from_secs(3600 * 2)),
1051 max_output_file_size: Some(ReadableSize::mb(7)),
1052 remote_compaction: false,
1053 fallback_to_local: true,
1054 }),
1055 compaction_override: false,
1056 storage: Some("S3".to_string()),
1057 append_mode: false,
1058 skip_wal: false,
1059 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1060 index_options: IndexOptions {
1061 inverted_index: InvertedIndexOptions {
1062 ignore_column_ids: vec![],
1063 segment_row_count: 512,
1064 },
1065 },
1066 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1067 merge_mode: Some(MergeMode::LastNonNull),
1068 sst_format: None,
1069 max_row_group_row_count: None,
1070 primary_key_encoding: None,
1071 write_buffer_size: None,
1072 };
1073 assert_eq!(options, got);
1074 }
1075
1076 #[test]
1077 fn test_max_row_group_row_count() {
1078 assert_eq!(None, RegionOptions::default().max_row_group_row_count);
1080 assert_eq!(
1081 DEFAULT_ROW_GROUP_SIZE,
1082 RegionOptions::default().row_group_size()
1083 );
1084
1085 let map = make_map(&[("max_row_group_row_count", "51200")]);
1087 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1088 assert_eq!(Some(51200), options.max_row_group_row_count);
1089 assert_eq!(51200, options.row_group_size());
1090
1091 let map = make_map(&[("max_row_group_row_count", "0")]);
1093 assert!(RegionOptions::try_from_options(RegionId::new(0, 0), &map).is_err());
1094 }
1095}