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::{
36 COMPACTION_OVERRIDE, FloatFieldEncoding, MAX_ROW_GROUP_ROW_COUNT_LIMIT,
37};
38use store_api::storage::{ColumnId, RegionId};
39use strum::EnumString;
40
41use crate::error::{InvalidRegionOptionsSnafu, JsonOptionsSnafu, Result};
42use crate::memtable::bulk::BulkMemtableConfig;
43use crate::sst::FormatType;
44use crate::sst::parquet::DEFAULT_ROW_GROUP_SIZE;
45
46const DEFAULT_INDEX_SEGMENT_ROW_COUNT: usize = 1024;
47const COMPACTION_TWCS_PREFIX: &str = "compaction.twcs.";
48const MEMTABLE_PARTITION_TREE_PREFIX: &str = "memtable.partition_tree.";
49const MEMTABLE_BULK_PREFIX: &str = "memtable.bulk.";
50const MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD: &str = "memtable.bulk.encode_bytes_threshold";
51
52const LEGACY_PARTITION_TREE_MEMTABLE_TYPE: &str = "partition_tree";
56
57pub(crate) fn parse_wal_options(
58 options_map: &HashMap<String, String>,
59) -> std::result::Result<WalOptions, serde_json::Error> {
60 options_map
61 .get(WAL_OPTIONS_KEY)
62 .map_or(Ok(WalOptions::default()), |encoded_wal_options| {
63 serde_json::from_str(encoded_wal_options)
64 })
65}
66
67#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, EnumString)]
69#[serde(rename_all = "snake_case")]
70#[strum(serialize_all = "snake_case")]
71pub enum MergeMode {
72 #[default]
74 LastRow,
75 LastNonNull,
77}
78
79#[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
85#[serde(default)]
86pub struct RegionOptions {
87 pub ttl: Option<TimeToLive>,
89 #[serde(with = "humantime_serde")]
92 pub auto_flush_interval: Option<Duration>,
93 pub compaction: CompactionOptions,
95 pub compaction_override: bool,
96 pub storage: Option<String>,
98 pub append_mode: bool,
100 pub skip_wal: bool,
102 pub wal_options: WalOptions,
104 pub index_options: IndexOptions,
106 pub memtable: Option<MemtableOptions>,
108 pub merge_mode: Option<MergeMode>,
111 pub sst_format: Option<FormatType>,
113 #[serde(skip_serializing_if = "Option::is_none")]
115 pub max_row_group_row_count: Option<usize>,
116 #[serde(skip_serializing_if = "Option::is_none")]
118 pub primary_key_encoding: Option<PrimaryKeyEncoding>,
119 #[serde(skip_serializing_if = "Option::is_none")]
123 pub write_buffer_size: Option<ReadableSize>,
124 #[serde(default, skip_serializing_if = "is_false")]
133 pub preserve_row_sequence: bool,
134 #[serde(default)]
136 pub float_field_encoding: FloatFieldEncoding,
137}
138
139fn is_false(value: &bool) -> bool {
140 !*value
141}
142
143impl RegionOptions {
144 pub fn validate(&self) -> Result<()> {
146 if self.append_mode {
147 ensure!(
148 self.merge_mode
149 .is_none_or(|mode| mode == MergeMode::LastRow),
150 InvalidRegionOptionsSnafu {
151 reason: "only last_row merge_mode is allowed when append_mode is enabled",
152 }
153 );
154 }
155 if let Some(auto_flush_interval) = self.auto_flush_interval {
156 ensure!(
157 auto_flush_interval > Duration::ZERO,
158 InvalidRegionOptionsSnafu {
159 reason: "auto_flush_interval must be greater than 0",
160 }
161 );
162 }
163 if let Some(row_count) = self.max_row_group_row_count {
164 ensure!(
165 row_count > 0 && row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT,
166 InvalidRegionOptionsSnafu {
167 reason: format!(
168 "max_row_group_row_count must be in (0, {MAX_ROW_GROUP_ROW_COUNT_LIMIT}], got {row_count}",
169 ),
170 }
171 );
172 }
173 if self.preserve_row_sequence {
174 ensure!(
175 self.append_mode,
176 InvalidRegionOptionsSnafu {
177 reason: "preserve_row_sequence is only supported for append-only tables (append_mode must be true)",
178 }
179 );
180 }
181 let CompactionOptions::Twcs(options) = &self.compaction;
182 ensure!(
183 options.active_window_l1_merge_trigger >= 2,
184 InvalidRegionOptionsSnafu {
185 reason: "active_window.l1_merge_trigger must be at least 2",
186 }
187 );
188 ensure!(
189 options.inactive_window_l1_merge_trigger >= 2,
190 InvalidRegionOptionsSnafu {
191 reason: "inactive_window.l1_merge_trigger must be at least 2",
192 }
193 );
194 Ok(())
195 }
196
197 pub fn row_group_size(&self) -> usize {
199 self.max_row_group_row_count
200 .unwrap_or(DEFAULT_ROW_GROUP_SIZE)
201 }
202
203 pub fn need_dedup(&self) -> bool {
205 !self.append_mode
206 }
207
208 pub fn merge_mode(&self) -> MergeMode {
210 self.merge_mode.unwrap_or_default()
211 }
212
213 pub fn auto_flush_interval_or(&self, default: Duration) -> Duration {
215 self.auto_flush_interval.unwrap_or(default)
216 }
217
218 pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
220 self.primary_key_encoding.unwrap_or_default()
221 }
222}
223
224impl RegionOptions {
225 pub fn try_from_options(
227 region_id: RegionId,
228 options_map: &HashMap<String, String>,
229 ) -> Result<Self> {
230 let value = options_map_to_value(options_map);
231 let json = serde_json::to_string(&value).context(JsonOptionsSnafu)?;
232
233 let options: RegionOptionsWithoutEnum =
237 serde_json::from_str(&json).context(JsonOptionsSnafu)?;
238 let has_compaction_type =
239 validate_enum_options(options_map, "compaction.type", &[COMPACTION_TWCS_PREFIX])?;
240 let compaction = if has_compaction_type {
241 serde_json::from_str(&json).context(JsonOptionsSnafu)?
242 } else {
243 CompactionOptions::default()
244 };
245
246 let wal_options = parse_wal_options(options_map).context(JsonOptionsSnafu)?;
247
248 let index_options: IndexOptions = serde_json::from_str(&json).context(JsonOptionsSnafu)?;
249 let is_legacy_partition_tree = options_map
250 .get("memtable.type")
251 .map(|s| s.eq_ignore_ascii_case(LEGACY_PARTITION_TREE_MEMTABLE_TYPE))
252 .unwrap_or(false);
253 let memtable = if validate_enum_options(
254 options_map,
255 "memtable.type",
256 &[MEMTABLE_PARTITION_TREE_PREFIX, MEMTABLE_BULK_PREFIX],
257 )? {
258 if is_legacy_partition_tree {
259 None
263 } else {
264 Some(serde_json::from_str(&json).context(JsonOptionsSnafu)?)
265 }
266 } else {
267 None
268 };
269
270 let mut sst_format = options.sst_format;
273 if is_legacy_partition_tree {
274 info!(
275 "Region {} specified the removed partition_tree memtable; \
276 overriding memtable to the default and SST format to flat",
277 region_id
278 );
279 sst_format = Some(FormatType::Flat);
280 }
281
282 if matches!(memtable, Some(MemtableOptions::Bulk(_))) {
285 if let Some(format) = sst_format
286 && format != FormatType::Flat
287 {
288 info!(
289 "Region {} uses bulk memtable; overriding sst_format from {:?} to flat",
290 region_id, format
291 );
292 }
293 sst_format = Some(FormatType::Flat);
294 }
295
296 let float_field_encoding = options.float_field_encoding.unwrap_or_default();
297
298 let compaction_override_flag = options_map
299 .get(COMPACTION_OVERRIDE)
300 .map(|v| matches!(v.to_lowercase().as_str(), "true" | "1"))
301 .unwrap_or(false);
302 let compaction_override = has_compaction_type || compaction_override_flag;
303 let primary_key_encoding = options_map
304 .get(PRIMARY_KEY_ENCODING)
305 .or_else(|| options_map.get(MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING))
306 .map(|v| match v.to_lowercase().as_str() {
307 "dense" => Ok(PrimaryKeyEncoding::Dense),
308 "sparse" => Ok(PrimaryKeyEncoding::Sparse),
309 _ => Err(InvalidRegionOptionsSnafu {
310 reason: format!("Invalid primary key encoding: {v}"),
311 }
312 .build()),
313 })
314 .transpose()?;
315
316 let opts = RegionOptions {
317 ttl: options.ttl,
318 auto_flush_interval: options.auto_flush_interval,
319 compaction,
320 compaction_override,
321 storage: options.storage,
322 append_mode: options.append_mode,
323 skip_wal: options.skip_wal,
324 wal_options,
325 index_options,
326 memtable,
327 merge_mode: options.merge_mode,
328 sst_format,
329 max_row_group_row_count: options.max_row_group_row_count,
330 primary_key_encoding,
331 write_buffer_size: options.write_buffer_size,
332 preserve_row_sequence: options.preserve_row_sequence,
333 float_field_encoding,
334 };
335 opts.validate()?;
336
337 Ok(opts)
338 }
339
340 pub(crate) fn try_from_options_with_bulk_config(
342 region_id: RegionId,
343 options_map: &HashMap<String, String>,
344 default_bulk_config: &BulkMemtableConfig,
345 ) -> Result<Self> {
346 let mut options = Self::try_from_options(region_id, options_map)?;
347 if !options_map.contains_key(MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD)
348 && let Some(MemtableOptions::Bulk(config)) = &mut options.memtable
349 {
350 config.encode_bytes_threshold = default_bulk_config.encode_bytes_threshold;
351 }
352 Ok(options)
353 }
354}
355
356#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
358#[serde(tag = "compaction.type")]
359#[serde(rename_all = "snake_case")]
360pub enum CompactionOptions {
361 #[serde(with = "prefix_twcs")]
363 Twcs(TwcsOptions),
364}
365
366impl CompactionOptions {
367 pub(crate) fn time_window(&self) -> Option<Duration> {
368 match self {
369 CompactionOptions::Twcs(opts) => opts.time_window,
370 }
371 }
372
373 pub(crate) fn remote_compaction(&self) -> bool {
374 match self {
375 CompactionOptions::Twcs(opts) => opts.remote_compaction,
376 }
377 }
378
379 pub(crate) fn fallback_to_local(&self) -> bool {
380 match self {
381 CompactionOptions::Twcs(opts) => opts.fallback_to_local,
382 }
383 }
384}
385
386impl Default for CompactionOptions {
387 fn default() -> Self {
388 Self::Twcs(TwcsOptions::default())
389 }
390}
391
392#[serde_as]
394#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
395#[serde(default)]
396pub struct TwcsOptions {
397 #[serde_as(as = "DisplayFromStr")]
399 #[serde(rename = "active_window.trigger_file_num", alias = "trigger_file_num")]
400 pub active_window_trigger_file_num: usize,
401 #[serde_as(as = "DisplayFromStr")]
403 #[serde(rename = "active_window.l1_merge_trigger")]
404 pub active_window_l1_merge_trigger: usize,
405 #[serde_as(as = "DisplayFromStr")]
407 #[serde(rename = "inactive_window.trigger_file_num")]
408 pub inactive_window_trigger_file_num: usize,
409 #[serde_as(as = "DisplayFromStr")]
411 #[serde(rename = "inactive_window.l1_merge_trigger")]
412 pub inactive_window_l1_merge_trigger: usize,
413 #[serde(with = "humantime_serde")]
415 pub time_window: Option<Duration>,
416 pub max_output_file_size: Option<ReadableSize>,
420 #[serde_as(as = "DisplayFromStr")]
422 pub remote_compaction: bool,
423 #[serde_as(as = "DisplayFromStr")]
425 pub fallback_to_local: bool,
426}
427
428with_prefix!(prefix_twcs "compaction.twcs.");
429
430impl TwcsOptions {
431 pub fn time_window_seconds(&self) -> Option<i64> {
433 self.time_window.and_then(|window| {
434 let window_secs = window.as_secs();
435 if window_secs == 0 {
436 None
437 } else {
438 window_secs.try_into().ok()
439 }
440 })
441 }
442}
443
444impl Default for TwcsOptions {
445 fn default() -> Self {
446 Self {
447 active_window_trigger_file_num: 4,
448 active_window_l1_merge_trigger: 16,
449 inactive_window_trigger_file_num: 2,
450 inactive_window_l1_merge_trigger: 8,
451 time_window: None,
452 max_output_file_size: Some(ReadableSize::mb(512)),
453 remote_compaction: false,
454 fallback_to_local: true,
455 }
456 }
457}
458
459#[serde_as]
462#[derive(Debug, Deserialize)]
463#[serde(default)]
464struct RegionOptionsWithoutEnum {
465 write_buffer_size: Option<ReadableSize>,
466 ttl: Option<TimeToLive>,
468 #[serde(with = "humantime_serde")]
469 auto_flush_interval: Option<Duration>,
470 storage: Option<String>,
471 #[serde_as(as = "DisplayFromStr")]
472 append_mode: bool,
473 #[serde_as(as = "DisplayFromStr")]
474 skip_wal: bool,
475 #[serde_as(as = "NoneAsEmptyString")]
476 merge_mode: Option<MergeMode>,
477 #[serde_as(as = "NoneAsEmptyString")]
478 sst_format: Option<FormatType>,
479 #[serde_as(as = "NoneAsEmptyString")]
480 max_row_group_row_count: Option<usize>,
481 #[serde_as(as = "DisplayFromStr")]
482 preserve_row_sequence: bool,
483 #[serde(rename = "experimental_sst_float_field_encoding")]
484 #[serde_as(as = "NoneAsEmptyString")]
485 float_field_encoding: Option<FloatFieldEncoding>,
486}
487
488impl Default for RegionOptionsWithoutEnum {
489 fn default() -> Self {
490 let options = RegionOptions::default();
491 RegionOptionsWithoutEnum {
492 write_buffer_size: options.write_buffer_size,
493 ttl: options.ttl,
494 auto_flush_interval: options.auto_flush_interval,
495 storage: options.storage,
496 append_mode: options.append_mode,
497 skip_wal: options.skip_wal,
498 merge_mode: options.merge_mode,
499 sst_format: options.sst_format,
500 max_row_group_row_count: options.max_row_group_row_count,
501 preserve_row_sequence: options.preserve_row_sequence,
502 float_field_encoding: Some(options.float_field_encoding),
503 }
504 }
505}
506
507with_prefix!(prefix_inverted_index "index.inverted_index.");
508
509#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
511#[serde(default)]
512pub struct IndexOptions {
513 #[serde(flatten, with = "prefix_inverted_index")]
515 pub inverted_index: InvertedIndexOptions,
516}
517
518#[serde_as]
520#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
521#[serde(default)]
522pub struct InvertedIndexOptions {
523 #[serde(deserialize_with = "deserialize_ignore_column_ids")]
526 #[serde(serialize_with = "serialize_ignore_column_ids")]
527 pub ignore_column_ids: Vec<ColumnId>,
528
529 #[serde_as(as = "DisplayFromStr")]
531 pub segment_row_count: usize,
532}
533
534impl Default for InvertedIndexOptions {
535 fn default() -> Self {
536 Self {
537 ignore_column_ids: Vec::new(),
538 segment_row_count: DEFAULT_INDEX_SEGMENT_ROW_COUNT,
539 }
540 }
541}
542
543#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
545#[serde(tag = "memtable.type", rename_all = "snake_case")]
546pub enum MemtableOptions {
547 TimeSeries,
548 #[serde(with = "prefix_bulk")]
549 Bulk(BulkMemtableConfig),
550}
551
552with_prefix!(prefix_bulk "memtable.bulk.");
553
554fn deserialize_ignore_column_ids<'de, D>(deserializer: D) -> Result<Vec<ColumnId>, D::Error>
555where
556 D: Deserializer<'de>,
557{
558 let s: String = Deserialize::deserialize(deserializer)?;
559 let mut column_ids = Vec::new();
560 if s.is_empty() {
561 return Ok(column_ids);
562 }
563 for item in s.split(',') {
564 let column_id = item.parse().map_err(D::Error::custom)?;
565 column_ids.push(column_id);
566 }
567 Ok(column_ids)
568}
569
570fn serialize_ignore_column_ids<S>(column_ids: &[ColumnId], serializer: S) -> Result<S::Ok, S::Error>
571where
572 S: serde::Serializer,
573{
574 let s = column_ids
575 .iter()
576 .map(|id| id.to_string())
577 .collect::<Vec<_>>()
578 .join(",");
579 serializer.serialize_str(&s)
580}
581
582fn options_map_to_value(options: &HashMap<String, String>) -> Value {
586 let map = options
587 .iter()
588 .map(|(key, value)| {
589 if value.eq_ignore_ascii_case("null") {
591 (key.clone(), Value::Null)
592 } else {
593 (key.clone(), Value::from(value.clone()))
594 }
595 })
596 .collect();
597 Value::Object(map)
598}
599
600fn validate_enum_options(
608 options_map: &HashMap<String, String>,
609 enum_tag_key: &str,
610 enum_option_prefixes: &[&str],
611) -> Result<bool> {
612 let mut has_enum_options = false;
613 let mut has_tag = false;
614 for key in options_map.keys() {
615 if key == enum_tag_key {
616 has_tag = true;
617 } else if !has_enum_options
618 && enum_option_prefixes
619 .iter()
620 .any(|prefix| key.starts_with(prefix))
621 {
622 has_enum_options = true;
623 }
624
625 if has_tag && has_enum_options {
626 break;
627 }
628 }
629
630 ensure!(
632 has_tag || !has_enum_options,
633 InvalidRegionOptionsSnafu {
634 reason: format!("missing key {} in options", enum_tag_key),
635 }
636 );
637
638 Ok(has_tag)
639}
640
641#[cfg(test)]
642mod tests {
643 use common_error::ext::ErrorExt;
644 use common_error::status_code::StatusCode;
645 use common_wal::options::KafkaWalOptions;
646 use store_api::mito_engine_options::{
647 EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, SKIP_WAL_KEY, WRITE_BUFFER_SIZE_KEY,
648 };
649
650 use super::*;
651
652 fn make_map(options: &[(&str, &str)]) -> HashMap<String, String> {
653 options
654 .iter()
655 .map(|(k, v)| (k.to_string(), v.to_string()))
656 .collect()
657 }
658
659 #[test]
660 fn test_empty_region_options() {
661 let map = make_map(&[]);
662 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
663 assert_eq!(RegionOptions::default(), options);
664 }
665
666 #[test]
667 fn test_float_field_encoding_defaults_and_parses() {
668 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &make_map(&[])).unwrap();
669 assert_eq!(FloatFieldEncoding::Default, options.float_field_encoding);
670
671 let map = make_map(&[(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, "default")]);
672 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
673 assert_eq!(FloatFieldEncoding::Default, options.float_field_encoding);
674
675 let map = make_map(&[(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, "byte_stream_split")]);
676 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
677 assert_eq!(
678 FloatFieldEncoding::ByteStreamSplit,
679 options.float_field_encoding
680 );
681
682 let map = make_map(&[(EXPERIMENTAL_SST_FLOAT_FIELD_ENCODING, "invalid")]);
683 assert!(RegionOptions::try_from_options(RegionId::new(0, 0), &map).is_err());
684 }
685
686 #[test]
687 fn test_with_ttl() {
688 let map = make_map(&[("ttl", "7d")]);
689 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
690 let expect = RegionOptions {
691 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
692 ..Default::default()
693 };
694 assert_eq!(expect, options);
695 }
696
697 #[test]
698 fn test_with_skip_wal() {
699 let map = make_map(&[(SKIP_WAL_KEY, "true")]);
700 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
701 assert!(options.skip_wal);
702 }
703
704 #[test]
705 fn test_with_auto_flush_interval() {
706 let map = make_map(&[("auto_flush_interval", "5m")]);
707 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
708 let expect = RegionOptions {
709 auto_flush_interval: Some(Duration::from_secs(5 * 60)),
710 ..Default::default()
711 };
712 assert_eq!(expect, options);
713 }
714
715 #[test]
716 fn test_with_zero_auto_flush_interval() {
717 let map = make_map(&[("auto_flush_interval", "0s")]);
718 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
719 assert!(
720 err.to_string().contains("auto_flush_interval"),
721 "unexpected error: {err}"
722 );
723 }
724
725 #[test]
726 fn test_with_write_buffer_size() {
727 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "128MiB")]);
728 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
729 let expect = RegionOptions {
730 write_buffer_size: Some(ReadableSize::mb(128)),
731 ..Default::default()
732 };
733 assert_eq!(expect, options);
734 }
735
736 #[test]
737 fn test_with_zero_write_buffer_size() {
738 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "0")]);
739 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
740 let expect = RegionOptions {
741 write_buffer_size: Some(ReadableSize::mb(0)),
742 ..Default::default()
743 };
744 assert_eq!(expect, options);
745 }
746
747 #[test]
748 fn test_with_storage() {
749 let map = make_map(&[("storage", "S3")]);
750 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
751 let expect = RegionOptions {
752 storage: Some("S3".to_string()),
753 ..Default::default()
754 };
755 assert_eq!(expect, options);
756 }
757
758 #[test]
759 fn test_without_compaction_type() {
760 let map = make_map(&[
761 ("compaction.twcs.trigger_file_num", "8"),
762 ("compaction.twcs.time_window", "2h"),
763 ]);
764 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
765 assert_eq!(StatusCode::InvalidArguments, err.status_code());
766 }
767
768 #[test]
769 fn test_with_compaction_type() {
770 let map = make_map(&[
771 ("compaction.twcs.active_window.trigger_file_num", "8"),
772 ("compaction.twcs.active_window.l1_merge_trigger", "16"),
773 ("compaction.twcs.inactive_window.trigger_file_num", "2"),
774 ("compaction.twcs.inactive_window.l1_merge_trigger", "12"),
775 ("compaction.twcs.time_window", "2h"),
776 ("compaction.type", "twcs"),
777 ]);
778 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
779 let expect = RegionOptions {
780 compaction: CompactionOptions::Twcs(TwcsOptions {
781 active_window_trigger_file_num: 8,
782 active_window_l1_merge_trigger: 16,
783 inactive_window_trigger_file_num: 2,
784 inactive_window_l1_merge_trigger: 12,
785 time_window: Some(Duration::from_secs(3600 * 2)),
786 ..Default::default()
787 }),
788 compaction_override: true,
789 ..Default::default()
790 };
791 assert_eq!(expect, options);
792 }
793
794 #[test]
795 fn test_twcs_window_trigger_below_two_is_accepted_for_compatibility() {
796 let map = make_map(&[
799 ("compaction.twcs.active_window.trigger_file_num", "1"),
800 ("compaction.type", "twcs"),
801 ]);
802 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
803 let CompactionOptions::Twcs(twcs) = &options.compaction;
804 assert_eq!(1, twcs.active_window_trigger_file_num);
805
806 let map = make_map(&[
807 ("compaction.twcs.inactive_window.trigger_file_num", "1"),
808 ("compaction.type", "twcs"),
809 ]);
810 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
811 let CompactionOptions::Twcs(twcs) = &options.compaction;
812 assert_eq!(1, twcs.inactive_window_trigger_file_num);
813 }
814
815 #[test]
816 fn test_active_window_l1_merge_trigger_below_two_is_rejected() {
817 let map = make_map(&[
818 ("compaction.twcs.active_window.l1_merge_trigger", "1"),
819 ("compaction.type", "twcs"),
820 ]);
821
822 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
823 assert_eq!(StatusCode::InvalidArguments, err.status_code());
824 }
825
826 #[test]
827 fn test_inactive_window_l1_merge_trigger_defaults_to_eight() {
828 let value = serde_json::to_value(TwcsOptions::default()).unwrap();
829 assert_eq!(
830 Some("8"),
831 value
832 .get("inactive_window.l1_merge_trigger")
833 .and_then(|value| value.as_str())
834 );
835 }
836
837 #[test]
838 fn test_inactive_window_l1_merge_trigger_below_two_is_rejected() {
839 let map = make_map(&[
840 ("compaction.twcs.inactive_window.l1_merge_trigger", "1"),
841 ("compaction.type", "twcs"),
842 ]);
843
844 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
845 assert_eq!(StatusCode::InvalidArguments, err.status_code());
846 }
847
848 #[test]
849 fn test_with_compaction_override_true_without_compaction_type() {
850 let map = make_map(&[(COMPACTION_OVERRIDE, "true")]);
851 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
852 let expect = RegionOptions {
853 compaction_override: true,
854 ..Default::default()
855 };
856 assert_eq!(expect, options);
857 }
858
859 #[test]
860 fn test_with_compaction_override_false_without_compaction_type() {
861 let map = make_map(&[(COMPACTION_OVERRIDE, "false")]);
862 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
863 assert_eq!(RegionOptions::default(), options);
864 }
865
866 #[test]
867 fn test_compaction_twcs_options_still_require_compaction_type_with_override() {
868 let map = make_map(&[
869 (COMPACTION_OVERRIDE, "true"),
870 ("compaction.twcs.time_window", "2h"),
871 ]);
872 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
873 assert_eq!(StatusCode::InvalidArguments, err.status_code());
874 }
875
876 fn test_with_wal_options(wal_options: &WalOptions) -> bool {
877 let encoded_wal_options = serde_json::to_string(&wal_options).unwrap();
878 let map = make_map(&[(WAL_OPTIONS_KEY, &encoded_wal_options)]);
879 let got = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
880 let expect = RegionOptions {
881 wal_options: wal_options.clone(),
882 ..Default::default()
883 };
884 expect == got
885 }
886
887 #[test]
888 fn test_with_index() {
889 let map = make_map(&[
890 ("index.inverted_index.ignore_column_ids", "1,2,3"),
891 ("index.inverted_index.segment_row_count", "512"),
892 ]);
893 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
894 let expect = RegionOptions {
895 index_options: IndexOptions {
896 inverted_index: InvertedIndexOptions {
897 ignore_column_ids: vec![1, 2, 3],
898 segment_row_count: 512,
899 },
900 },
901 ..Default::default()
902 };
903 assert_eq!(expect, options);
904 }
905
906 #[test]
908 fn test_with_any_wal_options() {
909 let all_wal_options = [
910 WalOptions::RaftEngine,
911 WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
912 ];
913 all_wal_options.iter().all(test_with_wal_options);
914 }
915
916 #[test]
917 fn test_with_memtable() {
918 let map = make_map(&[("memtable.type", "time_series")]);
919 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
920 let expect = RegionOptions {
921 memtable: Some(MemtableOptions::TimeSeries),
922 ..Default::default()
923 };
924 assert_eq!(expect, options);
925
926 let map = make_map(&[("memtable.type", "bulk")]);
927 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
928 let expect = RegionOptions {
929 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
930 sst_format: Some(FormatType::Flat),
931 ..Default::default()
932 };
933 assert_eq!(expect, options);
934
935 let map = make_map(&[
936 ("memtable.type", "bulk"),
937 ("memtable.bulk.merge_threshold", "7"),
938 ("memtable.bulk.encode_row_threshold", "11"),
939 ("memtable.bulk.encode_bytes_threshold", "13"),
940 ("memtable.bulk.max_merge_groups", "17"),
941 ]);
942 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
943 let expect = RegionOptions {
944 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
945 merge_threshold: 7,
946 encode_row_threshold: 11,
947 encode_bytes_threshold: 13,
948 max_merge_groups: 17,
949 })),
950 sst_format: Some(FormatType::Flat),
951 ..Default::default()
952 };
953 assert_eq!(expect, options);
954
955 let map = make_map(&[("memtable.type", "partition_tree")]);
958 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
959 let expect = RegionOptions {
960 memtable: None,
961 sst_format: Some(FormatType::Flat),
962 ..Default::default()
963 };
964 assert_eq!(expect, options);
965
966 let map = make_map(&[
968 ("memtable.type", "partition_tree"),
969 ("memtable.partition_tree.index_max_keys_per_shard", "2048"),
970 ("memtable.partition_tree.fork_dictionary_bytes", "128M"),
971 ]);
972 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
973 let expect = RegionOptions {
974 memtable: None,
975 sst_format: Some(FormatType::Flat),
976 ..Default::default()
977 };
978 assert_eq!(expect, options);
979 }
980
981 #[test]
982 fn test_primary_key_encoding() {
983 let map = make_map(&[("primary_key_encoding", "sparse")]);
985 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
986 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
987 assert_eq!(
988 options.primary_key_encoding,
989 Some(PrimaryKeyEncoding::Sparse)
990 );
991
992 let map = make_map(&[
994 ("memtable.type", "partition_tree"),
995 ("memtable.partition_tree.primary_key_encoding", "sparse"),
996 ]);
997 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
998 assert_eq!(options.memtable, None);
999 assert_eq!(options.sst_format, Some(FormatType::Flat));
1000 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
1001
1002 let map = make_map(&[("primary_key_encoding", "bogus")]);
1004 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1005 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1006 }
1007
1008 #[test]
1009 fn test_legacy_partition_tree_overrides_sst_format() {
1010 let map = make_map(&[
1013 ("memtable.type", "partition_tree"),
1014 ("sst_format", "primary_key"),
1015 ]);
1016 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
1017 assert_eq!(options.memtable, None);
1018 assert_eq!(options.sst_format, Some(FormatType::Flat));
1019 }
1020
1021 #[test]
1022 fn test_bulk_memtable_overrides_sst_format() {
1023 let map = make_map(&[("memtable.type", "bulk"), ("sst_format", "primary_key")]);
1027 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
1028 assert_eq!(
1029 options.memtable,
1030 Some(MemtableOptions::Bulk(BulkMemtableConfig::default()))
1031 );
1032 assert_eq!(options.sst_format, Some(FormatType::Flat));
1033 }
1034
1035 #[test]
1036 fn test_unknown_memtable_type() {
1037 let map = make_map(&[("memtable.type", "no_such_memtable")]);
1038 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1039 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1040 }
1041
1042 #[test]
1043 fn test_without_memtable_type() {
1044 let map = make_map(&[("memtable.partition_tree.index_max_keys_per_shard", "2048")]);
1045 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1046 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1047
1048 let map = make_map(&[("memtable.bulk.merge_threshold", "7")]);
1049 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1050 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1051 }
1052
1053 #[test]
1054 fn test_with_merge_mode() {
1055 let map = make_map(&[("merge_mode", "last_row")]);
1056 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1057 assert_eq!(MergeMode::LastRow, options.merge_mode());
1058
1059 let map = make_map(&[("merge_mode", "last_non_null")]);
1060 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1061 assert_eq!(MergeMode::LastNonNull, options.merge_mode());
1062
1063 let map = make_map(&[("merge_mode", "unknown")]);
1064 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1065 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1066 }
1067
1068 #[test]
1069 fn test_append_mode_allows_last_row_merge_mode() {
1070 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_row")]);
1071 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1072 assert!(options.append_mode);
1073 assert_eq!(MergeMode::LastRow, options.merge_mode());
1074
1075 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_non_null")]);
1076 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1077 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1078 }
1079
1080 #[test]
1081 fn test_with_all() {
1082 let wal_options = WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string()));
1083 let map = make_map(&[
1084 ("ttl", "7d"),
1085 ("compaction.twcs.trigger_file_num", "8"),
1086 ("compaction.twcs.max_output_file_size", "1GB"),
1087 ("compaction.twcs.time_window", "2h"),
1088 ("compaction.type", "twcs"),
1089 ("compaction.twcs.remote_compaction", "false"),
1090 ("compaction.twcs.fallback_to_local", "true"),
1091 ("storage", "S3"),
1092 ("append_mode", "false"),
1093 ("index.inverted_index.ignore_column_ids", "1,2,3"),
1094 ("index.inverted_index.segment_row_count", "512"),
1095 (
1096 WAL_OPTIONS_KEY,
1097 &serde_json::to_string(&wal_options).unwrap(),
1098 ),
1099 ("memtable.type", "bulk"),
1100 ("memtable.bulk.merge_threshold", "7"),
1101 ("memtable.bulk.encode_row_threshold", "11"),
1102 ("memtable.bulk.encode_bytes_threshold", "13"),
1103 ("memtable.bulk.max_merge_groups", "17"),
1104 ("merge_mode", "last_non_null"),
1105 ]);
1106 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1107 let expect = RegionOptions {
1108 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1109 auto_flush_interval: None,
1110 compaction: CompactionOptions::Twcs(TwcsOptions {
1111 active_window_trigger_file_num: 8,
1112 active_window_l1_merge_trigger: 16,
1113 inactive_window_trigger_file_num: 2,
1114 inactive_window_l1_merge_trigger: 8,
1115 time_window: Some(Duration::from_secs(3600 * 2)),
1116 max_output_file_size: Some(ReadableSize::gb(1)),
1117 remote_compaction: false,
1118 fallback_to_local: true,
1119 }),
1120 compaction_override: true,
1121 storage: Some("S3".to_string()),
1122 append_mode: false,
1123 skip_wal: false,
1124 wal_options,
1125 index_options: IndexOptions {
1126 inverted_index: InvertedIndexOptions {
1127 ignore_column_ids: vec![1, 2, 3],
1128 segment_row_count: 512,
1129 },
1130 },
1131 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
1132 merge_threshold: 7,
1133 encode_row_threshold: 11,
1134 encode_bytes_threshold: 13,
1135 max_merge_groups: 17,
1136 })),
1137 merge_mode: Some(MergeMode::LastNonNull),
1138 sst_format: Some(FormatType::Flat),
1139 max_row_group_row_count: None,
1140 primary_key_encoding: None,
1141 write_buffer_size: None,
1142 preserve_row_sequence: false,
1143 float_field_encoding: FloatFieldEncoding::default(),
1144 };
1145 assert_eq!(expect, options);
1146 }
1147
1148 #[test]
1149 fn test_with_preserve_row_sequence() {
1150 let map = make_map(&[("append_mode", "false"), ("preserve_row_sequence", "true")]);
1151 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1152 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1153 assert!(err.to_string().contains("preserve_row_sequence"));
1154
1155 let map = make_map(&[("append_mode", "true"), ("preserve_row_sequence", "true")]);
1156 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1157 assert!(options.append_mode);
1158 assert!(options.preserve_row_sequence);
1159
1160 let map = make_map(&[("append_mode", "true")]);
1161 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1162 assert!(!options.preserve_row_sequence);
1163
1164 let map = make_map(&[]);
1165 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1166 assert_eq!(RegionOptions::default(), options);
1167 }
1168
1169 #[test]
1170 fn test_region_options_serde() {
1171 let options = RegionOptions {
1172 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1173 auto_flush_interval: None,
1174 compaction: CompactionOptions::Twcs(TwcsOptions {
1175 active_window_trigger_file_num: 8,
1176 active_window_l1_merge_trigger: 8,
1177 inactive_window_trigger_file_num: 2,
1178 inactive_window_l1_merge_trigger: 8,
1179 time_window: Some(Duration::from_secs(3600 * 2)),
1180 max_output_file_size: None,
1181 remote_compaction: false,
1182 fallback_to_local: true,
1183 }),
1184 compaction_override: false,
1185 storage: Some("S3".to_string()),
1186 append_mode: false,
1187 skip_wal: false,
1188 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1189 index_options: IndexOptions {
1190 inverted_index: InvertedIndexOptions {
1191 ignore_column_ids: vec![1, 2, 3],
1192 segment_row_count: 512,
1193 },
1194 },
1195 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1196 merge_mode: Some(MergeMode::LastNonNull),
1197 sst_format: None,
1198 max_row_group_row_count: None,
1199 primary_key_encoding: None,
1200 write_buffer_size: Some(ReadableSize::mb(128)),
1201 preserve_row_sequence: true,
1202 float_field_encoding: FloatFieldEncoding::default(),
1203 };
1204 let region_options_json_str = serde_json::to_string(&options).unwrap();
1205 assert!(region_options_json_str.contains("preserve_row_sequence"));
1206 let got: RegionOptions = serde_json::from_str(®ion_options_json_str).unwrap();
1207 assert_eq!(options, got);
1208
1209 let old_region_options_json_str = r#"{"ttl":null}"#;
1211 let got: RegionOptions = serde_json::from_str(old_region_options_json_str).unwrap();
1212 assert_eq!(None, got.write_buffer_size);
1213 assert!(!got.preserve_row_sequence);
1214 assert_eq!(FloatFieldEncoding::Default, got.float_field_encoding);
1215 let CompactionOptions::Twcs(twcs) = got.compaction;
1216 assert_eq!(16, twcs.active_window_l1_merge_trigger);
1217 assert_eq!(8, twcs.inactive_window_l1_merge_trigger);
1218
1219 let default_json = serde_json::to_value(RegionOptions::default()).unwrap();
1220 assert!(default_json.get(WRITE_BUFFER_SIZE_KEY).is_none());
1221 }
1222
1223 #[test]
1224 fn test_region_options_str_serde() {
1225 let region_options_json_str = r#"{
1227 "ttl": "7days",
1228 "compaction": {
1229 "compaction.type": "twcs",
1230 "compaction.twcs.trigger_file_num": "8",
1231 "compaction.twcs.max_output_file_size": "7MB",
1232 "compaction.twcs.time_window": "2h"
1233 },
1234 "storage": "S3",
1235 "append_mode": false,
1236 "wal_options": {
1237 "wal.provider": "kafka",
1238 "wal.kafka.topic": "test_topic"
1239 },
1240 "index_options": {
1241 "index.inverted_index.ignore_column_ids": "",
1242 "index.inverted_index.segment_row_count": "512"
1243 },
1244 "memtable": {
1245 "memtable.type": "bulk"
1246 },
1247 "merge_mode": "last_non_null"
1248}"#;
1249 let got: RegionOptions = serde_json::from_str(region_options_json_str).unwrap();
1250 let options = RegionOptions {
1251 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1252 auto_flush_interval: None,
1253 compaction: CompactionOptions::Twcs(TwcsOptions {
1254 active_window_trigger_file_num: 8,
1255 active_window_l1_merge_trigger: 16,
1256 inactive_window_trigger_file_num: 2,
1257 inactive_window_l1_merge_trigger: 8,
1258 time_window: Some(Duration::from_secs(3600 * 2)),
1259 max_output_file_size: Some(ReadableSize::mb(7)),
1260 remote_compaction: false,
1261 fallback_to_local: true,
1262 }),
1263 compaction_override: false,
1264 storage: Some("S3".to_string()),
1265 append_mode: false,
1266 skip_wal: false,
1267 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1268 index_options: IndexOptions {
1269 inverted_index: InvertedIndexOptions {
1270 ignore_column_ids: vec![],
1271 segment_row_count: 512,
1272 },
1273 },
1274 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1275 merge_mode: Some(MergeMode::LastNonNull),
1276 sst_format: None,
1277 max_row_group_row_count: None,
1278 primary_key_encoding: None,
1279 write_buffer_size: None,
1280 preserve_row_sequence: false,
1281 float_field_encoding: FloatFieldEncoding::default(),
1282 };
1283 assert_eq!(options, got);
1284 }
1285
1286 #[test]
1287 fn test_max_row_group_row_count() {
1288 assert_eq!(None, RegionOptions::default().max_row_group_row_count);
1290 assert_eq!(
1291 DEFAULT_ROW_GROUP_SIZE,
1292 RegionOptions::default().row_group_size()
1293 );
1294
1295 let map = make_map(&[("max_row_group_row_count", "51200")]);
1297 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1298 assert_eq!(Some(51200), options.max_row_group_row_count);
1299 assert_eq!(51200, options.row_group_size());
1300
1301 let map = make_map(&[("max_row_group_row_count", "0")]);
1303 assert!(RegionOptions::try_from_options(RegionId::new(0, 0), &map).is_err());
1304 }
1305}