@@ -20,6 +20,7 @@ use std::collections::{HashMap, HashSet};
2020const DELETION_VECTORS_ENABLED_OPTION : & str = "deletion-vectors.enabled" ;
2121const DATA_EVOLUTION_ENABLED_OPTION : & str = "data-evolution.enabled" ;
2222const GLOBAL_INDEX_ENABLED_OPTION : & str = "global-index.enabled" ;
23+ const GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION : & str = "global-index.row-count-per-shard" ;
2324const SOURCE_SPLIT_TARGET_SIZE_OPTION : & str = "source.split.target-size" ;
2425const SOURCE_SPLIT_OPEN_FILE_COST_OPTION : & str = "source.split.open-file-cost" ;
2526const PARTITION_DEFAULT_NAME_OPTION : & str = "partition.default-name" ;
@@ -64,6 +65,7 @@ const DEFAULT_TARGET_FILE_SIZE: i64 = 256 * 1024 * 1024;
6465const DEFAULT_WRITE_PARQUET_BUFFER_SIZE : i64 = 256 * 1024 * 1024 ;
6566const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION : & str = "dynamic-bucket.target-row-num" ;
6667const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM : i64 = 200_000 ;
68+ const DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD : i64 = 100_000 ;
6769const BLOB_AS_DESCRIPTOR_OPTION : & str = "blob-as-descriptor" ;
6870const BLOB_DESCRIPTOR_FIELD_OPTION : & str = "blob-descriptor-field" ;
6971
@@ -222,6 +224,22 @@ impl<'a> CoreOptions<'a> {
222224 . unwrap_or ( false )
223225 }
224226
227+ pub fn global_index_row_count_per_shard ( & self ) -> crate :: Result < i64 > {
228+ let value = self
229+ . parse_i64_option ( GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION ) ?
230+ . unwrap_or ( DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD ) ;
231+ if value <= 0 {
232+ return Err ( crate :: Error :: DataInvalid {
233+ message : format ! (
234+ "Option '{}' must be greater than 0, got: {}" ,
235+ GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION , value
236+ ) ,
237+ source : None ,
238+ } ) ;
239+ }
240+ Ok ( value)
241+ }
242+
225243 pub fn source_split_target_size ( & self ) -> i64 {
226244 self . options
227245 . get ( SOURCE_SPLIT_TARGET_SIZE_OPTION )
@@ -536,6 +554,10 @@ mod tests {
536554
537555 assert_eq ! ( core_options. source_split_target_size( ) , 128 * 1024 * 1024 ) ;
538556 assert_eq ! ( core_options. source_split_open_file_cost( ) , 4 * 1024 * 1024 ) ;
557+ assert_eq ! (
558+ core_options. global_index_row_count_per_shard( ) . unwrap( ) ,
559+ 100_000
560+ ) ;
539561 }
540562
541563 #[ test]
@@ -549,11 +571,36 @@ mod tests {
549571 SOURCE_SPLIT_OPEN_FILE_COST_OPTION . to_string ( ) ,
550572 "8 mb" . to_string ( ) ,
551573 ) ,
574+ (
575+ GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION . to_string ( ) ,
576+ "2048" . to_string ( ) ,
577+ ) ,
552578 ] ) ;
553579 let core_options = CoreOptions :: new ( & options) ;
554580
555581 assert_eq ! ( core_options. source_split_target_size( ) , 256 * 1024 * 1024 ) ;
556582 assert_eq ! ( core_options. source_split_open_file_cost( ) , 8 * 1024 * 1024 ) ;
583+ assert_eq ! (
584+ core_options. global_index_row_count_per_shard( ) . unwrap( ) ,
585+ 2048
586+ ) ;
587+ }
588+
589+ #[ test]
590+ fn test_global_index_row_count_per_shard_rejects_invalid_values ( ) {
591+ for value in [ "0" , "-1" , "abc" ] {
592+ let options = HashMap :: from ( [ (
593+ GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION . to_string ( ) ,
594+ value. to_string ( ) ,
595+ ) ] ) ;
596+ let core = CoreOptions :: new ( & options) ;
597+
598+ let err = core
599+ . global_index_row_count_per_shard ( )
600+ . expect_err ( "invalid rows-per-shard should fail" ) ;
601+ assert ! ( matches!( err, crate :: Error :: DataInvalid { message, .. }
602+ if message. contains( GLOBAL_INDEX_ROW_COUNT_PER_SHARD_OPTION ) ) ) ;
603+ }
557604 }
558605
559606 #[ test]
0 commit comments