FE Configuration - Statistics and Storage
Follow these steps to configure Coordinator Node parameters in the PhoenixAI Cloud console:
- Sign in to the PhoenixAI Cloud console.
- On the Clusters page, click the cluster that you want to configure.
- On the cluster details page, click the Cluster parameters tab.
- In the Coordinator Node static configuration section, click View all parameter list.
- In the dialog box that appears, click New parameter.
- Enter the name of the parameter you want to configure in the Parameter key field, and the parameter value you want to set in the Value field.
note
PhoenixAI provides validation checks on parameter keys and value types. Only valid parameter keys and values can be applied. For the details of the parameters you can configure, see Usage notes.
- Click the save button next to the Value field to save the record of the parameter change.
- After configuring all the parameters you want to change, click Save changes to save the changes.
- To allow the changes to take effect, you can either manually suspend and resume the cluster whenever you find suitable, or click Apply to all nodes in the Coordinator Node static configuration section to restart the cluster instantly.
View FE configuration items
After your FE is started, you can run the ADMIN SHOW FRONTEND CONFIG command on your MySQL client to check the parameter configurations. If you want to query the configuration of a specific parameter, run the following command:
ADMIN SHOW FRONTEND CONFIG [LIKE "pattern"];
For detailed description of the returned fields, see ADMIN SHOW CONFIG.
You must have administrator privileges to run cluster administration-related commands.
Configure FE parameters
Configure FE dynamic parameters
You can configure or modify the settings of FE dynamic parameters using ADMIN SET FRONTEND CONFIG.
ADMIN SET FRONTEND CONFIG ("key" = "value");
The configuration changes made with ADMIN SET FRONTEND will be restored to the default values after the FE restarts. Therefore, we recommend that you also modify the configuration items in PhoenixAI Cloud console if you want the changes to be permanent.
Configure FE static parameters
Static parameters of an FE are set by changing them in the PhoenixAI Cloud console and restarting the FE to allow the changes to take effect.
This topic introduces the following types of FE configurations:
Statistic report
enable_collect_warehouse_metrics
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: When this item is set to
true, the system will collect and export per-warehouse metrics. Enabling it adds warehouse-level metrics (slot/usage/availability) to the metric output and increases metric cardinality and collection overhead. Disable it to omit warehouse-specific metrics and reduce CPU/network and monitoring storage cost. - Introduced in: v3.5.0
enable_http_detail_metrics
- Default: false
- Type: boolean
- Unit: -
- Is mutable: Yes
- Description: When true, the HTTP server computes and exposes detailed HTTP worker metrics (notably the
HTTP_WORKER_PENDING_TASKS_NUMgauge). Enabling this causes the server to iterate over Netty worker executors and callpendingTasks()on eachNioEventLoopto sum pending task counts; when disabled the gauge returns 0 to avoid that cost. This extra collection can be CPU- and latency-sensitive — enable only for debugging or detailed investigation. - Introduced in: v3.2.3
proc_profile_collect_time_s
- Default: 120
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: Duration in seconds for a single process profile collection. When
proc_profile_cpu_enableorproc_profile_mem_enableis set totrue, AsyncProfiler is started, the collector thread sleeps for this duration, then the profiler is stopped and the profile is written. Larger values increase sample coverage and file size but prolong profiler runtime and delay subsequent collections; smaller values reduce overhead but may produce insufficient samples. Ensure this value aligns with retention settings such asproc_profile_file_retained_daysandproc_profile_file_retained_size_bytes. - Introduced in: v3.2.12
low_cardinality_dict_cache_max_bytes
- Default: 1073741824
- Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Maximum total size (in bytes) of the low-cardinality global dictionary cache (
CacheDictManager). The cache is bounded by the combined byte size of its cached dictionaries rather than by entry count, so its memory footprint is bounded directly (each dictionary can be up to ~1 MB). When the limit is reached the least-valuable dictionaries are evicted, and affected columns fall back to non-dictionary query plans until re-collected. Changes apply to the live cache within one config-refresh cycle. The current tracked size is exported via thelow_cardinality_dict_cache_bytesmetric. - Introduced in: v4.1.0
enable_external_predicate_columns_collection
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to record predicate column usage (columns used in WHERE/JOIN/GROUP BY) for external (non-native) tables during query optimization. StarRocks uses this usage information to narrow down which columns ANALYZE collects statistics for on wide external tables. When disabled, external table predicate columns are not recorded, and ANALYZE falls back to collecting statistics for all columns.
- Introduced in: v4.2.0
statistic_external_predicate_columns_ttl_hours
- Default: 168
- Type: Long
- Unit: Hours
- Is mutable: Yes
- Description: The time-to-live (TTL) of recorded external table predicate column usage. Entries whose
last_usedtimestamp is older than this value are removed by the periodic vacuum job. Set to a negative value (e.g. -1) to disable vacuum. Defaults to a week because external table ANALYZE runs far less frequently than for internal tables, so a short TTL (matching the internal table's 24-hour default) would evict usage information between two collections. - Introduced in: v4.2.0
statistic_external_predicate_columns_cache_ttl_sec
- Default: 300
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The TTL of the in-memory cache that serves external table predicate column queries (for example, during automatic ANALYZE column selection). A shorter value makes newly recorded usage visible sooner but increases the query load on the underlying storage table; a longer value reduces that load at the cost of staleness.
- Introduced in: v4.2.0
Storage
allow_implicit_key_column_in_agg_add_column
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether
ALTER TABLE ... ADD COLUMNon an Aggregate table may create a key column when the new column specifies neither an aggregate function nor theKEYkeyword. Such a statement is ambiguous, and creating a key column changes the table's aggregation key and rewrites existing data. When set totrue, the column is created as a key column, which is the behavior in earlier versions. Set tofalseto reject the statement instead, so that the error names both options. This item is mutable but is not persisted across a restart unless it is set withWITH PERSISTENT. - Introduced in: v4.2.0
alter_table_timeout_second
- Default: 86400
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: The timeout duration for the schema change operation (ALTER TABLE).
- Introduced in: -
enable_concurrent_add_partition_during_alter
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: When
true, partition creation (manualALTER TABLE ... ADD PARTITION, automatic creation during loading, and the dynamic partition scheduler) is allowed to proceed concurrently with metadata-only alter operations that are provably safe — currently the shared-data ADD/DROP INDEX fast-path jobs and the transientUPDATING_METAstate of fast schema evolution — instead of rejecting the DDL or cancelling the alter job. Set tofalseto restore the legacy exclusive behavior. This setting only relaxes partition creation; all other alter jobs and all non-ADD PARTITIONDDL keep the legacy state checks. - Introduced in: -
capacity_used_percent_high_water
- Default: 0.75
- Type: double
- Unit: Fraction (0.0–1.0)
- Is mutable: Yes
- Description: The high-water threshold of disk capacity used percent (fraction of total capacity) used when computing backend load scores.
BackendLoadStatistic.calcSoreusescapacity_used_percent_high_waterto setLoadScore.capacityCoefficient: if a backend's used percent less than 0.5 the coefficient equal to 0.5; if used percent>capacity_used_percent_high_waterthe coefficient = 1.0; otherwise the coefficient transitions linearly with used percent via (2 * usedPercent - 0.5). When the coefficient is 1.0, the load score is driven entirely by capacity proportion; lower values increase the weight of replica count. Adjusting this value changes how aggressively the balancer penalizes backends with high disk utilization. - Introduced in: v3.2.0
catalog_trash_expire_second
- Default: 86400
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The longest duration the metadata can be retained after a database, table, or partition is dropped. If this duration expires, the data will be deleted and cannot be recovered through the RECOVER command.
- Introduced in: -
catalog_recycle_bin_erase_min_latency_ms
- Default: 600000
- Type: Long
- Unit: Milliseconds
- Is mutable: Yes
- Description: The minimum delay in milliseconds before the metadata is erased when a database, table, or partition is dropped. This avoids the erase log being written ahead of the drop log.
- Introduced in: -
catalog_recycle_bin_erase_max_operations_per_cycle
- Default: 500
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of erase operations per cycle for actually deleting databases, tables, or partitions from the recycle bin. The erase operation holds a lock, so one batch should not be too large.
- Introduced in: -
catalog_recycle_bin_erase_fail_retry_interval_ms
- Default: 60000
- Type: Long
- Unit: Milliseconds
- Is mutable: Yes
- Description: The retry interval in milliseconds when an erase operation in the recycle bin fails.
- Introduced in: -
check_consistency_default_timeout_second
- Default: 600
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The timeout duration for a replica consistency check. You can set this parameter based on the size of your tablet.
- Introduced in: -
consistency_check_cooldown_time_second
- Default: 24 * 3600
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: Controls the minimal interval (in seconds) required between consistency checks of the same tablet. During tablet selection, a tablet is considered eligible only if
tablet.getLastCheckTime()is less than(currentTimeMillis - consistency_check_cooldown_time_second * 1000). The default value (24 * 3600) enforces roughly one check per tablet per day to reduce backend disk I/O. Lowering this value increases check frequency and resource usage; raising it reduces I/O at the cost of slower detection of inconsistencies. The value is applied globally when filtering cooldowned tablets from an index's tablet list. - Introduced in: v3.5.5
consistency_check_end_time
- Default: "4"
- Type: String
- Unit: Hour of day (0-23)
- Is mutable: No
- Description: Specifies the end hour (hour-of-day) of the ConsistencyChecker work window. The value is parsed with SimpleDateFormat("HH") in the system time zone and accepted as 0–23 (single or two-digit). StarRocks uses it with
consistency_check_start_timeto decide when to schedule and add consistency-check jobs. Whenconsistency_check_start_timeis greater thanconsistency_check_end_time, the window spans midnight (for example, default isconsistency_check_start_time= "23" toconsistency_check_end_time= "4"). Whenconsistency_check_start_timeis equal toconsistency_check_end_time, the checker never runs. Parsing failure will cause FE startup to log an error and exit, so provide a valid hour string. - Introduced in: v3.2.0
consistency_check_start_time
- Default: "23"
- Type: String
- Unit: Hour of day (00-23)
- Is mutable: No
- Description: Specifies the start hour (hour-of-day) of the ConsistencyChecker work window. The value is parsed with SimpleDateFormat("HH") in the system time zone and accepted as 0–23 (single or two-digit). StarRocks uses it with
consistency_check_end_timeto decide when to schedule and add consistency-check jobs. Whenconsistency_check_start_timeis greater thanconsistency_check_end_time, the window spans midnight (for example, default isconsistency_check_start_time= "23" toconsistency_check_end_time= "4"). Whenconsistency_check_start_timeis equal toconsistency_check_end_time, the checker never runs. Parsing failure will cause FE startup to log an error and exit, so provide a valid hour string. - Introduced in: v3.2.0
consistency_tablet_meta_check_interval_ms
- Default: 2 * 3600 * 1000
- Type: Int
- Unit: Milliseconds
- Is mutable: Yes
- Description: Interval used by the ConsistencyChecker to run a full tablet-meta consistency scan between
TabletInvertedIndexandLocalMetastore. The daemon inrunAfterCatalogReadytriggers checkTabletMetaConsistency whencurrent time - lastTabletMetaCheckTimeexceeds this value. When an invalid tablet is first detected, itstoBeCleanedTimeis set tonow + (consistency_tablet_meta_check_interval_ms / 2)so actual deletion is delayed until a subsequent scan. Increase this value to reduce scan frequency and load (slower cleanup); decrease it to detect and remove stale tablets faster (higher overhead). - Introduced in: v3.2.0
default_replication_num
- Default: 3
- Type: Short
- Unit: -
- Is mutable: Yes
- Description: Sets the default number of replicas for each data partition when creating a table in StarRocks. This setting can be overridden when creating a table by specifying
replication_num=xin the CREATE TABLE DDL. - Introduced in: -
enable_auto_tablet_distribution
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to automatically set the number of buckets.
- If this parameter is set to
TRUE, you don't need to specify the number of buckets when you create a table or add a partition. StarRocks automatically determines the number of buckets. - If this parameter is set to
FALSE, you need to manually specify the number of buckets when you create a table or add a partition. If you do not specify the bucket count when adding a new partition to a table, the new partition inherits the bucket count set at the creation of the table. However, you can also manually specify the number of buckets for the new partition.
- If this parameter is set to
- Introduced in: v2.5.7
enable_experimental_rowstore
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable the hybrid row-column storage feature.
- Introduced in: v3.2.3
enable_fast_schema_evolution
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable Fast Schema Evolution for all tables within the StarRocks cluster. Valid values are
TRUEandFALSE(default). Enabling Fast Schema Evolution can increase the speed of schema changes and reduce resource usage when columns are added or dropped. - Introduced in: v3.2.0
NOTE
- StarRocks shared-data clusters supports this parameter from v3.3.0.
- If you need to configure the Fast Schema Evolution for a specific table, such as disabling Fast Schema Evolution for a specific table, you can set the table property
fast_schema_evolutionat table creation.
enable_online_optimize_table
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Controls whether StarRocks will use the non-blocking online optimization path when creating an optimize job. When
enable_online_optimize_tableis true and the target table meets compatibility checks (no partition/keys/sort specification, distribution is notRandomDistributionDesc, storage type is notCOLUMN_WITH_ROW, replicated storage enabled, and the table is not a cloud-native table or materialized view), the planner creates anOnlineOptimizeJobV2to perform optimization without blocking writes. If false or any compatibility condition fails, StarRocks falls back toOptimizeJobV2, which may block write operations during optimization. - Introduced in: v3.3.3, v3.4.0, v3.5.0
enable_strict_storage_medium_check
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether the FE strictly checks the storage medium of BEs when users create tables. If this parameter is set to
TRUE, the FE checks the storage medium of BEs when users create tables and returns an error if the storage medium of the BE is different from thestorage_mediumparameter specified in the CREATE TABLE statement. For example, the storage medium specified in the CREATE TABLE statement is SSD but the actual storage medium of BEs is HDD. As a result, the table creation fails. If this parameter isFALSE, the FE does not check the storage medium of BEs when users create a table. - Introduced in: -
max_bucket_number_per_partition
- Default: 1024
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of buckets can be created in a partition.
- Introduced in: v3.3.2
max_column_number_per_table
- Default: 10000
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of columns can be created in a table.
- Introduced in: v3.3.2
max_dynamic_partition_num
- Default: 500
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Limits the maximum number of partitions that can be created at once when analyzing or creating a dynamic-partitioned table. During dynamic partition property validation, the
systemtask_runs_max_history_numbercomputes expected partitions (end offset + history partition number) and throws a DDL error if that total exceedsmax_dynamic_partition_num. Raise this value only when you expect legitimately large partition ranges; increasing it allows more partitions to be created but can increase metadata size, scheduling work, and operational complexity. - Introduced in: v3.2.0
max_partition_number_per_table
- Default: 100000
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of partitions can be created in a table.
- Introduced in: v3.3.2
max_task_consecutive_fail_count
- Default: 10
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Maximum number of consecutive failures a task may have before the scheduler automatically suspends it. When
TaskSource.MV.equals(task.getSource())andmax_task_consecutive_fail_countare greater than 0, if a task's consecutive failure counter reaches or exceedsmax_task_consecutive_fail_count, the task is suspended via the TaskManager and, for materialized-view tasks, the materialized view is inactivated. An exception is thrown indicating suspension and how to reactivate (for example,ALTER MATERIALIZED VIEW <mv_name> ACTIVE). Set this item to 0 or a negative value to disable automatic suspension. - Introduced in: -
partition_recycle_retention_period_secs
- Default: 1800
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The metadata retention time for the partition that is dropped by INSERT OVERWRITE or materialized view refresh operations. Note that such metadata cannot be recovered by executing RECOVER.
- Introduced in: v3.5.9
recover_with_empty_tablet
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to replace a lost or corrupted tablet replica with an empty one. If a tablet replica is lost or corrupted, data queries on this tablet or other healthy tablets may fail. Replacing the lost or corrupted tablet replica with an empty tablet ensures that the query can still be executed. However, the result may be incorrect because data is lost. The default value is
FALSE, which means lost or corrupted tablet replicas are not replaced with empty ones, and the query fails. - Introduced in: -
storage_usage_hard_limit_percent
- Default: 95
- Alias:
storage_flood_stage_usage_percent - Type: Int
- Unit: -
- Is mutable: Yes
- Description: Hard limit of the storage usage percentage in a BE directory. If the storage usage (in percentage) of the BE storage directory exceeds this value and the remaining storage space is less than
storage_usage_hard_limit_reserve_bytes, Load and Restore jobs are rejected. You need to set this item together with the BE configuration itemstorage_flood_stage_usage_percentto allow the configurations to take effect. - Introduced in: -
storage_usage_hard_limit_reserve_bytes
- Default: 100 * 1024 * 1024 * 1024
- Alias:
storage_flood_stage_left_capacity_bytes - Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Hard limit of the remaining storage space in a BE directory. If the remaining storage space in the BE storage directory is less than this value and the storage usage (in percentage) exceeds
storage_usage_hard_limit_percent, Load and Restore jobs are rejected. You need to set this item together with the BE configuration itemstorage_flood_stage_left_capacity_bytesto allow the configurations to take effect. - Introduced in: -
storage_usage_soft_limit_percent
- Default: 90
- Alias:
storage_high_watermark_usage_percent - Type: Int
- Unit: -
- Is mutable: Yes
- Description: Soft limit of the storage usage percentage in a BE directory. If the storage usage (in percentage) of the BE storage directory exceeds this value and the remaining storage space is less than
storage_usage_soft_limit_reserve_bytes, tablets cannot be cloned into this directory. - Introduced in: -
storage_usage_soft_limit_reserve_bytes
- Default: 200 * 1024 * 1024 * 1024
- Alias:
storage_min_left_capacity_bytes - Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Soft limit of the remaining storage space in a BE directory. If the remaining storage space in the BE storage directory is less than this value and the storage usage (in percentage) exceeds
storage_usage_soft_limit_percent, tablets cannot be cloned into this directory. - Introduced in: -
tablet_checker_lock_time_per_cycle_ms
- Default: 1000
- Type: Int
- Unit: Milliseconds
- Is mutable: Yes
- Description: The maximum lock hold time per cycle for tablet checker before releasing and reacquiring the table lock. Values less than 100 will be treated as 100.
- Introduced in: v3.5.9, v4.0.2
tablet_create_timeout_second
- Default: 10
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: The timeout duration for creating a tablet. The default value is changed from 1 to 10 from v3.1 onwards.
- Introduced in: -
tablet_delete_timeout_second
- Default: 2
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: The timeout duration for deleting a tablet.
- Introduced in: -
tablet_sched_balance_load_disk_safe_threshold
- Default: 0.5
- Alias:
balance_load_disk_safe_threshold - Type: Double
- Unit: -
- Is mutable: Yes
- Description: The percentage threshold for determining whether the disk usage of BEs is balanced. If the disk usage of all BEs is lower than this value, it is considered balanced. If the disk usage is greater than this value and the difference between the highest and lowest BE disk usage is greater than 10%, the disk usage is considered unbalanced and a tablet re-balancing is triggered.
- Introduced in: -
tablet_sched_balance_load_score_threshold
- Default: 0.1
- Alias:
balance_load_score_threshold - Type: Double
- Unit: -
- Is mutable: Yes
- Description: The percentage threshold for determining whether the load of a BE is balanced. If a BE has a lower load than the average load of all BEs and the difference is greater than this value, this BE is in a low load state. On the contrary, if a BE has a higher load than the average load and the difference is greater than this value, this BE is in a high load state.
- Introduced in: -
tablet_sched_be_down_tolerate_time_s
- Default: 900
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The maximum duration the scheduler allows for a BE node to remain inactive. After the time threshold is reached, tablets on that BE node will be migrated to other active BE nodes.
- Introduced in: v2.5.7
tablet_sched_disable_balance
- Default: false
- Alias:
disable_balance - Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to disable tablet balancing.
TRUEindicates that tablet balancing is disabled.FALSEindicates that tablet balancing is enabled. - Introduced in: -
tablet_sched_disable_colocate_balance
- Default: false
- Alias:
disable_colocate_balance - Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to disable replica balancing for Colocate Table.
TRUEindicates replica balancing is disabled.FALSEindicates replica balancing is enabled. - Introduced in: -
tablet_sched_max_balancing_tablets
- Default: 500
- Alias:
max_balancing_tablets - Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of tablets that can be balanced at the same time. If this value is exceeded, tablet re-balancing will be skipped.
- Introduced in: -
tablet_sched_max_clone_task_timeout_sec
- Default: 2 * 60 * 60
- Alias:
max_clone_task_timeout_sec - Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description:The maximum timeout duration for cloning a tablet.
- Introduced in: -
tablet_sched_max_not_being_scheduled_interval_ms
- Default: 15 * 60 * 1000
- Type: Long
- Unit: Milliseconds
- Is mutable: Yes
- Description: When the tablet clone tasks are being scheduled, if a tablet has not been scheduled for the specified time in this parameter, StarRocks gives it a higher priority to schedule it as soon as possible.
- Introduced in: -
tablet_sched_max_scheduling_tablets
- Default: 10000
- Alias:
max_scheduling_tablets - Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of tablets that can be scheduled at the same time. If the value is exceeded, tablet balancing and repair checks will be skipped.
- Introduced in: -
tablet_sched_min_clone_task_timeout_sec
- Default: 3 * 60
- Alias:
min_clone_task_timeout_sec - Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The minimum timeout duration for cloning a tablet.
- Introduced in: -
tablet_sched_num_based_balance_threshold_ratio
- Default: 0.5
- Alias: -
- Type: Double
- Unit: -
- Is mutable: Yes
- Description: Doing num based balance may break the disk size balance, but the maximum gap between disks cannot exceed
tablet_sched_num_based_balance_threshold_ratio*tablet_sched_balance_load_score_threshold. If there are tablets in the cluster that are constantly balancing from A to B and B to A, reduce this value. If you want the tablet distribution to be more balanced, increase this value. - Introduced in: - 3.1
tablet_sched_repair_delay_factor_second
- Default: 60
- Alias:
tablet_repair_delay_factor_second - Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The interval at which replicas are repaired, in seconds.
- Introduced in: -
tablet_sched_slot_num_per_path
- Default: 8
- Alias:
schedule_slot_num_per_path - Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of tablet-related tasks that can run concurrently in a BE storage directory. From v2.5 onwards, the default value of this parameter is changed from
4to8. - Introduced in: -
tablet_sched_storage_cooldown_second
- Default: -1
- Alias:
storage_cooldown_second - Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: The latency of automatic cooling starting from the time of table creation. The default value
-1specifies that automatic cooling is disabled. If you want to enable automatic cooling, set this parameter to a value greater than-1. - Introduced in: -
tablet_stat_update_interval_second
- Default: 300
- Type: Int
- Unit: Seconds
- Is mutable: Yes
- Description: The time interval at which the FE retrieves tablet statistics from each BE.
- Introduced in: -
enable_range_distribution
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to use the Range-based Distribution semantic as the default table distribution when a table is created without a
DISTRIBUTED BYclause. This configuration only takes effect in shared-data mode; it has no effect in shared-nothing mode. While it isfalse, such a table uses the previous default distribution behavior instead (a PRIMARY KEY table defaults to hash, a DUPLICATE KEY table to random, and an AGGREGATE or UNIQUE KEY table requires an explicitDISTRIBUTED BYclause); set it totrueto make the Range-based Distribution semantic the default. A materialized view created without aDISTRIBUTED BYclause additionally requiresenable_mv_range_distribution. - Introduced in: v4.1.0
enable_mv_range_distribution
- Default: false
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to use the Range-based Distribution semantic as the default distribution of an asynchronous materialized view that is created without a
DISTRIBUTED BYclause. Tables are not affected by this configuration. The default selects the Range-based Distribution semantic only when this configuration andenable_range_distributionare bothtrue, in shared-data mode. Otherwise the materialized view uses the previous default distribution behavior (a materialized view that is maintained incrementally defaults to hash over its key columns, and any other materialized view to random), even where a table would be range-distributed. - Introduced in: v26.2
tablet_reshard_max_parallel_tablets
- Default: 10240
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of tablets that can be split or merged in parallel.
- Introduced in: v4.1.0
tablet_reshard_target_size
- Default: 10737418240 (10 GB)
- Type: Int
- Unit: Bytes
- Is mutable: Yes
- Description: The target size of the tablets after the SPLIT or MERGE operation.
- Introduced in: v4.1.0
tablet_reshard_max_split_count
- Default: 1024
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of new tablets that an old tablet can be split into.
- Introduced in: v4.1.0
tablet_reshard_orderby_max_split_count
- Default: 2
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of new tablets one source tablet may be split into when the split drags a full UNSHARE rewrite behind it, that is, on a range-distributed PRIMARY KEY table whose
ORDER BYkey differs from its primary key. Such a split cannot range-filter the parent's shared segments, so every child is rewritten wholesale and a wide fan-out multiplies that read amplification. Further clamped bytablet_reshard_max_split_count. Values less than or equal to1disable this extra clamp. - Introduced in: -
tablet_reshard_orderby_max_split_tablets_per_job
- Default: 0
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of source tablets one split job may split when the split drags a full UNSHARE rewrite behind it. The largest tablets are chosen first. Note that this bounds the split fan-out, not the rewrite itself: every untouched sibling still becomes an identical tablet in the replacement index, and the UNSHARE compaction is partition-wide, so those are rewritten as well. Values less than or equal to
0mean the compute-node count of the warehouse. - Introduced in: -
tablet_reshard_orderby_split_interval_second
- Default: 180
- Type: Int
- Unit: Second
- Is mutable: Yes
- Description: The quiet period after the previous tablet reshard job on a table finishes, before automatic splitting may trigger again, for tables whose split drags a full UNSHARE rewrite behind it. It gives size-tiered compaction a window to drain the small files that accumulated while the partition's compaction slot was held. Values less than or equal to
0disable the wait. Note that the interval can only be enforced while the previous job is still retained, that is, up totablet_reshard_history_job_keep_max_ms. - Introduced in: -
tablet_reshard_min_split_size
- Default: 2147483648 (2 GB)
- Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: The minimum size of a tablet produced by tablet pre-split. It bounds compute-node alignment during pre-split so that a small load on a large cluster is not split into many tiny tablets. It is also the smallest target size automatic splitting will aim at: while a materialized index holds fewer tablets than its warehouse has compute nodes (capped by
tablet_reshard_max_split_count), splitting aims at the size that would give it one tablet per such slot, floored at this value, so a tablet splits once it is worth at least two of that target. Raising this value therefore also delays that splitting, and setting it at or abovetablet_reshard_target_sizeturns it off, leaving only the size-based rule. Should be no larger thantablet_reshard_target_size. - Introduced in: v4.1.0
tablet_reshard_history_job_max_keep_ms
- Default: 259200000 (72 hours)
- Type: Int
- Unit: Milliseconds
- Is mutable: Yes
- Description: The maximum retention time of historical tablet SPLIT/MERGE jobs.
- Introduced in: v4.1.0
tablet_reshard_colocate_checker_membership_batch_size
- Default: 1000
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of tablets that the range-colocate checker sends to StarOS in a single
getShardInfomembership-read batch RPC on shared-data clusters. Values less than1are treated as1. - Introduced in: v4.1.3
tablet_reshard_colocate_checker_convergence_batch_size
- Default: 64
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: The maximum number of PACK shard groups that the range-colocate checker sends to StarOS in a single
queryShardGroupStableplacement-convergence batch RPC. Each group's stability check is computed server-side, so a smaller batch bounds per-RPC latency, and the full result is assembled across repeated calls. Values less than1are treated as1. - Introduced in: v4.1.3
tablet_reshard_colocate_checker_convergence_cache_ttl_ms
- Default: 1000
- Type: Long
- Unit: Milliseconds
- Is mutable: Yes
- Description: The TTL of the range-colocate checker's placement-convergence negative cache. Within this window, a PACK shard group that StarOS last reported as not yet converged is not re-queried, which throttles the per-tick
queryShardGroupStableload while the group is still migrating. Only not-converged results are cached, so a stale entry can only delay the group's flip to stable by up to this window, never cause a premature flip. Values less than or equal to0disable the cache. - Introduced in: v4.1.3
enable_tablet_pre_split_for_insert_from_files
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable Sample-Based Tablet Pre-Split for
INSERT INTO ... SELECT FROM FILES()loads. On by default as of v4.1.0. Set tofalseto disable cluster-wide. The session variableenable_tablet_pre_splitmust also betruefor pre-split to run. - Introduced in: v4.1.0
enable_tablet_pre_split_for_broker_load
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable Sample-Based Tablet Pre-Split for Broker Load. On by default as of v4.1.0. Set to
falseto disable cluster-wide. The session variableenable_tablet_pre_splitmust also betruefor pre-split to run. - Introduced in: v4.1.0
enable_tablet_pre_split_for_insert_from_table
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable Sample-Based Tablet Pre-Split for
INSERT INTO ... SELECT FROM <table>loads whose source is an internal OLAP table or a table in an external catalog, such as Hive, Iceberg, Paimon, Delta Lake, Hudi, JDBC, or Elasticsearch. Views are not supported as a source. When an external source cannot be sized from its table statistics, the load skips pre-split. The feature supports automatic range-partition targets, including explicitly named real or temporary partitions and both static and dynamicINSERT OVERWRITE. ForINSERT INTO, it also supports manually range-partitioned targets, whose existing empty partitions are split without creating any partition. On by default as of v4.1.0. Set tofalseto disable cluster-wide. The session variableenable_tablet_pre_splitmust also betruefor pre-split to run. To roll back, set tofalse; new INSERT-from-table loads will skip pre-split immediately. - Introduced in: v4.1.0
enable_tablet_pre_split_for_mv_refresh
- Default: true
- Type: Boolean
- Unit: -
- Is mutable: Yes
- Description: Whether to enable Sample-Based Tablet Pre-Split for the refresh of a range-distributed incremental materialized view. Such a view is keyed by a hidden row-id column whose value domain is known in advance, so its boundaries are derived rather than sampled and no data is read. Set to
falseto disable cluster-wide. The session variableenable_tablet_pre_splitmust also betruefor pre-split to run. - Introduced in: v4.2.0
tablet_pre_split_pre_submit_timeout_seconds
- Default: 300
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: Wall-clock budget for the pre-submit phase of Sample-Based Tablet Pre-Split (sample + plan boundaries + build reshard job). On expiry the coordinator skips pre-split and the load proceeds against the original single tablet. Default 300s: the data-tier sampler can take tens of seconds on large datasets / slow object storage (a ~40GB many-file Parquet load sampled in ~78s in testing), and this budget mainly bites large loads — exactly the ones pre-split benefits; small loads sample in well under a second regardless. The load stays
PENDINGfor at most this long during sampling, so keep it below the load's own timeout. - Introduced in: v4.1.0
tablet_pre_split_post_submit_wait_seconds
- Default: 300
- Type: Long
- Unit: Seconds
- Is mutable: Yes
- Description: Maximum time the coordinator will wait for an admitted Sample-Based Tablet Pre-Split reshard job to reach
FINISHED. Both INSERT-from-FILES and Broker Load synchronously wait and on expiry proceed without abort — the load then plans against the currently visible tablet layout (still the original layout if the daemon hasn't transitioned, or partially / fully post-split if the daemon raced past the wait); thetablet_pre_split_post_submit_hard_capcounter records the timeout. The strictrunPreSplitwrapper used by tests aborts the calling load viaPreSplitPostSubmitTimeoutException. For Broker Load the wait runs after the broker pending task resolves file statuses but beforebeginTxnopensT_load— it occupies apending_load_task_schedulerthread for at most this many seconds per table, so sizemax_broker_load_job_concurrencyaccordingly when many concurrent Broker Loads target a pre-splittable layout. Operator note: the Broker Load remainsPENDINGinSHOW LOADduring the wait and is still subject to its owntimeoutSecond— set this well below the smallest Broker Load timeout in normal use. - Introduced in: v4.1.0
tablet_pre_split_sample_byte_limit
- Default: 16777216 (16 MiB)
- Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Soft byte cap on the FE-side accumulation buffer of the data-tier reservoir sampler used by Sample-Based Tablet Pre-Split. The sampler stops reading once accumulated values exceed this limit. The first row is always admitted so an oversize row still produces a non-empty sample.
- Introduced in: v4.1.0
tablet_pre_split_data_tier_scan_byte_limit
- Default: 4294967296 (4 GiB)
- Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Soft limit on the source-file bytes the data tier of Sample-Based Tablet Pre-Split scans for a Broker Load or
INSERT INTO ... SELECT FROM FILES()load. When the input is larger, the sample reads a subset of the files instead of every file, so sampling time no longer grows with the input. The exception is a partition column read from the file data rather than from the path, or path and literal partition columns mixed: every file is then still scanned, because a subset could miss whole partitions. Files are sorted by path and picked at even byte intervals, which favors larger files. When the partition column comes from the file path (COLUMNS FROM PATHorcolumns_from_path), each partition gets a share of the limit in proportion to its bytes and at least one file (unless the load has more partitions thantablet_pre_split_data_tier_max_scan_files), and each partition's size is taken from the bytes of all of its files. Files are taken whole, so the scan can exceed the limit, for example when a single file is larger than it. The tablet count is still sized from the whole input. Set to0to scan every file.
tablet_pre_split_data_tier_min_scan_files
- Default: 64
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Minimum number of files the data tier of Sample-Based Tablet Pre-Split scans once it samples a subset of a load's files (see
tablet_pre_split_data_tier_scan_byte_limit), so that the sample spans enough independent files even when each file is large or holds only a narrow range of the sort key. When the partition column comes from the file path, the minimum is split across partitions in proportion to their bytes. Meeting it can take the scan abovetablet_pre_split_data_tier_scan_byte_limit, but no files are added for it once the scan reaches four times that limit, so large files do not multiply the scan time. A positivetablet_pre_split_data_tier_max_scan_filestakes precedence over this minimum.
tablet_pre_split_data_tier_max_scan_files
- Default: 512
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Maximum number of files the data tier of Sample-Based Tablet Pre-Split scans once it samples a subset of a load's files (see
tablet_pre_split_data_tier_scan_byte_limit). Every file in the subset costs one metadata lookup before any data is read, so without a cap a subset of many small files could spend minutes on lookups alone. When the partition column comes from the file path, the cap is split across partitions in proportion to their bytes; every partition still keeps at least one file, so the cap can be exceeded by up to one file per partition. When there are more partitions than the cap, only the heaviest partitions are sampled, one file each. The cap takes precedence overtablet_pre_split_data_tier_min_scan_files. Set to0or a negative value to remove the cap.
tablet_pre_split_meta_tier_overlap_threshold
- Default: 0.3
- Type: Double
- Unit: -
- Is mutable: Yes
- Description: Maximum overlap fraction tolerated when Sample-Based Tablet Pre-Split's meta tier (Parquet/ORC row-group metadata) computes boundaries. Above this threshold the cumulative-row count stops being monotone in sorted-min order so meta tier falls back to data tier (row sampling).
- Introduced in: v4.1.0
tablet_pre_split_meta_tier_footer_read_parallelism
- Default: 16
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Number of Parquet/ORC footers the Sample-Based Tablet Pre-Split meta tier reads concurrently from a file-backed load source (
FILES()or Broker Load). Footer reads are independent per file and the sampler sorts the aggregated statistics, so concurrency only cuts the wall time of the pre-split hook (each footer is a remote round-trip; a many-file source otherwise serializes hundreds of round-trips). Set to1to disable concurrency. - Introduced in: v4.1.0
tablet_pre_split_max_partitions_per_load
- Default: 32
- Type: Int
- Unit: -
- Is mutable: Yes
- Description: Maximum number of predicted target partitions a single Sample-Based Tablet Pre-Split invocation will operate on. Excess predicted partitions (the lightest ones: fewest input bytes when the data tier sampled a subset of files whose partition values come from the file paths, otherwise fewest sampled rows) are dropped and fall back to runtime auto-create with no pre-split. Bounds hook latency on pathological multi-partition loads. Set to zero or a negative value to disable the cap.
- Introduced in: v4.1.0
tablet_pre_split_target_size
- Default: 0
- Type: Long
- Unit: Bytes
- Is mutable: Yes
- Description: Target tablet size that Sample-Based Tablet Pre-Split sizes a load's split count against.
0(the default) inheritstablet_reshard_target_size. Lower it to give a load more write parallelism without shrinking every tablet in the cluster: the background tablet split/merge daemon keeps measuring againsttablet_reshard_target_size, so it merges the finer tablets back together after the load finishes. This matters most for a load that writes brand-new range-distributed partitions (for example the replacement partitions of anINSERT OVERWRITE), which start from a single catch-all tablet and would otherwise be written by a single backend.
Rolling back Sample-Based Tablet Pre-Split
To disable the feature safely before a downgrade or during a production rollback:
-
Set all four pre-split flags to
false:enable_tablet_pre_split_for_insert_from_files,enable_tablet_pre_split_for_broker_load,enable_tablet_pre_split_for_insert_from_table, andenable_tablet_pre_split_for_mv_refresh. New loads will skip pre-split immediately. -
Wait for in-flight reshard jobs created by pre-split to drain. Monitor them with the following query:
SELECT DB_NAME, TABLE_NAME, JOB_TYPE, JOB_STATEFROM information_schema.tablet_reshard_jobsWHERE JOB_STATE NOT IN ('FINISHED', 'ABORTED');The rollback is complete once this query returns no rows. A job in any non-final state, including
PENDING,PREPARING,RUNNING,CLEANING, andABORTING, is still in flight; onlyFINISHEDandABORTEDare final. -
Proceed with the downgrade. The substrate (External-Boundaries Tablet Split) remains available regardless of the pre-split feature flag.
Behavioral notes for multi-partition Sample-Based Tablet Pre-Split (P2-a)
The multi-partition path extends Sample-Based Tablet Pre-Split to loads that target many partitions in one statement. Three operational caveats apply:
- Broker Load triggering-load asymmetry. The multi-partition pre-split hook fires from
BrokerLoadJob.createLoadingTaskaftertask.prepare()has built the load's sink plan against the catalog as it existed at that moment. For Broker Load, pre-created partitions and the post-reshard tablet layout are therefore only visible to subsequent loads on the same table — the triggering Broker Load itself runs against the original layout and uses BE runtime auto-create for any partitions it touches. INSERT-from-FILES (where the hook fires beforeStatementPlanner.plan()) is unaffected and benefits in the same load. - Pre-created partition leak on subsequent INSERT failure. When pre-create succeeds and the triggering INSERT later fails for unrelated reasons (FILES schema mismatch, BE crash, load timeout, etc.), the empty pre-created partitions remain in the catalog. This matches the semantics of
ALTER TABLE ADD PARTITION, which also leaves a partition behind on subsequent failure. Operators who care can drop the empty partitions manually withALTER TABLE ... DROP PARTITION; in practice empty partitions are cheap and the next retry of the load will reuse them. - Manually range-partitioned targets. For a table with user-declared RANGE partitions, pre-split never creates a partition. Each sampled row is mapped onto the existing partition whose range contains it, and a value outside every declared range is dropped from the plan. Only partitions that are still empty and hold a single tablet are split, typically a partition just added with
ALTER TABLE ... ADD PARTITION; a load whose target partitions all hold data already skips pre-split without sampling the source.INSERT OVERWRITE, LIST partitioning, and expression partitioning are not covered for manually partitioned tables.
Production deployment guidance
Set enable_execute_script_on_frontend = false in production. Sample-Based Tablet Pre-Split exposes no SQL surface that depends on FE-side script execution; the production code paths sample through the connector + planner directly. Leaving enable_execute_script_on_frontend = true widens the FE attack surface without enabling any pre-split functionality, so the safe default for production clusters is to keep it off.