risingwave_common/config/
streaming.rs1use std::time::Duration;
16
17use risingwave_common_proc_macro::serde_prefix_all;
18
19use super::*;
20
21mod async_stack_trace;
22mod cache_refill;
23mod join_encoding_type;
24mod over_window;
25
26pub use async_stack_trace::*;
27pub use cache_refill::*;
28pub use join_encoding_type::*;
29pub use over_window::*;
30
31#[serde_with::apply(Option => #[serde(with = "none_as_empty_string")])]
33#[derive(Clone, Debug, Serialize, Deserialize, DefaultFromSerde, ConfigDoc)]
34pub struct StreamingConfig {
35 #[serde(default = "default::streaming::in_flight_barrier_nums")]
38 pub in_flight_barrier_nums: usize,
39
40 #[serde(default = "default::streaming::snapshot_backfill_finish_max_lagged_barriers")]
43 pub snapshot_backfill_finish_max_lagged_barriers: usize,
44
45 #[serde(default = "default::streaming::snapshot_backfill_barrier_amplification_factor")]
48 pub snapshot_backfill_barrier_amplification_factor: usize,
49
50 #[serde(default)]
53 pub actor_runtime_worker_threads_num: Option<usize>,
54
55 #[serde(default = "default::streaming::async_stack_trace")]
57 pub async_stack_trace: AsyncStackTraceOption,
58
59 #[serde(default)]
60 #[config_doc(nested)]
61 pub developer: StreamingDeveloperConfig,
62
63 #[serde(default = "default::streaming::unique_user_stream_errors")]
65 pub unique_user_stream_errors: usize,
66
67 #[serde(default = "default::streaming::unsafe_disable_strict_consistency")]
69 pub unsafe_disable_strict_consistency: bool,
70
71 #[serde(default, flatten)]
72 #[config_doc(omitted)]
73 pub unrecognized: Unrecognized<Self>,
74}
75
76#[serde_prefix_all("stream_", mode = "alias")]
80#[serde_with::apply(Option => #[serde(with = "none_as_empty_string")])]
81#[derive(Clone, Debug, Serialize, Deserialize, DefaultFromSerde, ConfigDoc)]
82pub struct StreamingDeveloperConfig {
83 #[serde(default = "default::developer::stream_enable_executor_row_count")]
87 pub enable_executor_row_count: bool,
88
89 #[serde(default = "default::developer::connector_message_buffer_size")]
92 pub connector_message_buffer_size: usize,
93
94 #[serde(default = "default::developer::unsafe_stream_extreme_cache_size")]
96 pub unsafe_extreme_cache_size: usize,
97
98 #[serde(default = "default::developer::stream_topn_cache_min_capacity")]
100 pub topn_cache_min_capacity: usize,
101
102 #[serde(default = "default::developer::stream_chunk_size")]
104 pub chunk_size: usize,
105
106 #[serde(default = "default::developer::stream_exchange_initial_permits")]
109 pub exchange_initial_permits: usize,
110
111 #[serde(default = "default::developer::stream_exchange_batched_permits")]
114 pub exchange_batched_permits: usize,
115
116 #[serde(default = "default::developer::stream_exchange_concurrent_barriers")]
118 pub exchange_concurrent_barriers: usize,
119
120 #[serde(default = "default::developer::stream_exchange_concurrent_dispatchers")]
125 pub exchange_concurrent_dispatchers: usize,
126
127 #[serde(default = "default::developer::stream_project_expr_concurrency")]
132 pub project_expr_concurrency: usize,
133
134 #[serde(default = "default::developer::stream_project_expr_inflight_request_concurrency")]
142 pub project_expr_inflight_request_concurrency: usize,
143
144 #[serde(default = "default::developer::stream_dml_channel_initial_permits")]
147 pub dml_channel_initial_permits: usize,
148
149 #[serde(default = "default::developer::stream_hash_agg_max_dirty_groups_heap_size")]
151 pub hash_agg_max_dirty_groups_heap_size: usize,
152
153 #[serde(default = "default::developer::memory_controller_threshold_aggressive")]
154 pub memory_controller_threshold_aggressive: f64,
155
156 #[serde(default = "default::developer::memory_controller_threshold_graceful")]
157 pub memory_controller_threshold_graceful: f64,
158
159 #[serde(default = "default::developer::memory_controller_threshold_stable")]
160 pub memory_controller_threshold_stable: f64,
161
162 #[serde(default = "default::developer::memory_controller_eviction_factor_aggressive")]
163 pub memory_controller_eviction_factor_aggressive: f64,
164
165 #[serde(default = "default::developer::memory_controller_eviction_factor_graceful")]
166 pub memory_controller_eviction_factor_graceful: f64,
167
168 #[serde(default = "default::developer::memory_controller_eviction_factor_stable")]
169 pub memory_controller_eviction_factor_stable: f64,
170
171 #[serde(default = "default::developer::memory_controller_update_interval_ms")]
172 pub memory_controller_update_interval_ms: usize,
173
174 #[serde(default = "default::developer::memory_controller_sequence_tls_step")]
175 pub memory_controller_sequence_tls_step: u64,
176
177 #[serde(default = "default::developer::memory_controller_sequence_tls_lag")]
178 pub memory_controller_sequence_tls_lag: u64,
179
180 #[serde(default = "default::developer::stream_enable_arrangement_backfill")]
181 #[deprecated(
184 note = "Deprecated and ignored for new streaming jobs. Arrangement backfill is always used as the fallback backfill type."
185 )]
186 pub enable_arrangement_backfill: bool,
187
188 #[serde(default = "default::developer::stream_enable_snapshot_backfill")]
189 pub enable_snapshot_backfill: bool,
194
195 #[serde(default = "default::developer::stream_high_join_amplification_threshold")]
196 pub high_join_amplification_threshold: usize,
199
200 #[serde(default = "default::developer::stream_high_gap_fill_amplification_threshold")]
201 pub high_gap_fill_amplification_threshold: usize,
204
205 #[serde(default = "default::developer::enable_actor_tokio_metrics")]
207 pub enable_actor_tokio_metrics: bool,
208
209 #[serde(default = "default::developer::stream_exchange_connection_pool_size")]
212 pub(super) exchange_connection_pool_size: Option<u16>,
213
214 #[serde(default = "default::developer::stream_enable_auto_schema_change")]
216 pub enable_auto_schema_change: bool,
217
218 #[serde(default = "default::developer::enable_shared_source")]
219 pub enable_shared_source: bool,
224
225 #[serde(default = "default::developer::switch_jdbc_pg_to_native")]
226 pub switch_jdbc_pg_to_native: bool,
229
230 #[serde(default = "default::developer::stream_max_barrier_batch_size")]
232 pub max_barrier_batch_size: u32,
233
234 #[serde(default = "default::developer::streaming_hash_join_entry_state_max_rows")]
237 pub hash_join_entry_state_max_rows: usize,
238
239 #[serde(default = "default::developer::streaming_join_hash_map_evict_interval_rows")]
242 pub join_hash_map_evict_interval_rows: u32,
243
244 #[serde(default = "default::developer::streaming_now_progress_ratio")]
245 pub now_progress_ratio: Option<f32>,
246
247 #[serde(default = "default::developer::enable_explain_analyze_stats")]
249 pub enable_explain_analyze_stats: bool,
250
251 #[serde(default)]
252 pub compute_client_config: RpcClientConfig,
253
254 #[serde(default = "default::developer::stream_snapshot_iter_rebuild_interval_secs")]
256 pub snapshot_iter_rebuild_interval_secs: u64,
257
258 #[serde(default = "default::developer::iceberg_list_interval_sec")]
260 pub iceberg_list_interval_sec: u64,
261
262 #[serde(default = "default::developer::iceberg_fetch_batch_size")]
264 pub iceberg_fetch_batch_size: u64,
265
266 #[serde(default = "default::developer::iceberg_sink_positional_delete_cache_size")]
268 pub iceberg_sink_positional_delete_cache_size: usize,
269
270 #[serde(default = "default::developer::iceberg_sink_write_parquet_max_row_group_rows")]
272 pub iceberg_sink_write_parquet_max_row_group_rows: usize,
273
274 #[serde(default = "default::developer::materialize_force_overwrite_on_no_check")]
277 pub materialize_force_overwrite_on_no_check: bool,
278
279 #[serde(default = "default::streaming::default_enable_mem_preload_state_table")]
282 pub default_enable_mem_preload_state_table: bool,
283
284 #[serde(default)]
287 pub mem_preload_state_table_ids_whitelist: Vec<u32>,
288
289 #[serde(default)]
292 pub mem_preload_state_table_ids_blacklist: Vec<u32>,
293
294 #[serde(default)]
297 pub aggressive_noop_update_elimination: bool,
298
299 #[serde(default = "default::developer::refresh_scheduler_interval_sec")]
301 pub refresh_scheduler_interval_sec: u64,
302
303 #[serde(default)]
305 pub join_encoding_type: JoinEncodingType,
306
307 #[serde(default = "default::developer::sync_log_store_pause_duration_ms")]
311 pub sync_log_store_pause_duration_ms: usize,
312
313 #[serde(default = "default::developer::sync_log_store_buffer_size")]
315 pub sync_log_store_buffer_size: usize,
316
317 #[serde(default = "default::developer::disable_sync_log_store_dispatcher")]
319 pub disable_sync_log_store_dispatcher: bool,
320
321 #[serde(default)]
324 pub over_window_cache_policy: OverWindowCachePolicy,
325
326 #[serde(default = "default::developer::enable_state_table_vnode_stats_pruning")]
332 pub enable_state_table_vnode_stats_pruning: bool,
333
334 #[serde(default = "default::developer::cache_refill_policy")]
337 pub cache_refill_policy: CacheRefillPolicy,
338
339 #[serde(default = "default::developer::enable_vnode_key_stats_for_materialize")]
341 pub enable_vnode_key_stats_for_materialize: bool,
342
343 #[serde(default = "default::developer::max_concurrent_kv_log_store_historical_read")]
348 pub max_concurrent_kv_log_store_historical_read: usize,
349
350 #[serde(default, flatten)]
351 #[serde_prefix_all(skip)]
352 #[config_doc(omitted)]
353 pub unrecognized: Unrecognized<Self>,
354}
355
356impl StreamingDeveloperConfig {
357 pub fn snapshot_iter_rebuild_interval(&self) -> Duration {
358 let rebuild_interval = if self.snapshot_iter_rebuild_interval_secs < 10 {
359 tracing::warn!(
360 "too small rebuild_interval {} second. rewrite to 10",
361 self.snapshot_iter_rebuild_interval_secs
362 );
363 10
364 } else {
365 self.snapshot_iter_rebuild_interval_secs
366 };
367 Duration::from_secs(rebuild_interval)
368 }
369}
370
371impl StreamingConfig {
372 pub fn unrecognized_keys(&self) -> impl Iterator<Item = String> {
374 std::iter::from_coroutine(
375 #[coroutine]
376 || {
377 for k in self.unrecognized.inner().keys() {
378 yield format!("streaming.{k}");
379 }
380 for k in self.developer.unrecognized.inner().keys() {
381 yield format!("streaming.developer.{k}");
382 }
383 },
384 )
385 }
386}
387
388pub mod default {
389 pub use crate::config::default::developer;
390
391 pub mod streaming {
392 use tracing::info;
393
394 use crate::config::AsyncStackTraceOption;
395 use crate::util::env_var::env_var_is_true;
396
397 pub fn in_flight_barrier_nums() -> usize {
398 10000
401 }
402
403 pub fn snapshot_backfill_finish_max_lagged_barriers() -> usize {
404 100
405 }
406
407 pub fn snapshot_backfill_barrier_amplification_factor() -> usize {
408 1
409 }
410
411 pub fn async_stack_trace() -> AsyncStackTraceOption {
412 AsyncStackTraceOption::default()
413 }
414
415 pub fn unique_user_stream_errors() -> usize {
416 10
417 }
418
419 pub fn unsafe_disable_strict_consistency() -> bool {
420 false
421 }
422
423 pub fn default_enable_mem_preload_state_table() -> bool {
424 if env_var_is_true("DEFAULT_ENABLE_MEM_PRELOAD_STATE_TABLE") {
425 info!("enabled mem_preload_state_table globally by env var");
426 true
427 } else {
428 false
429 }
430 }
431 }
432}