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