1mod block_stream;
16mod compaction_executor;
17mod compaction_filter;
18pub mod compaction_utils;
19mod iceberg_compaction;
20use risingwave_hummock_sdk::compact_task::{CompactTask, ValidationTask};
21use risingwave_pb::compactor::{DispatchCompactionTaskRequest, dispatch_compaction_task_request};
22use risingwave_pb::hummock::PbCompactTask;
23use risingwave_pb::hummock::report_compaction_task_request::{
24 Event as ReportCompactionTaskEvent, HeartBeat as SharedHeartBeat,
25 ReportTask as ReportSharedTask,
26};
27use risingwave_pb::iceberg_compaction::{
28 SubscribeIcebergCompactionEventRequest, SubscribeIcebergCompactionEventResponse,
29 subscribe_iceberg_compaction_event_request,
30};
31use risingwave_rpc_client::GrpcCompactorProxyClient;
32use thiserror_ext::AsReport;
33use tokio::sync::mpsc;
34use tonic::Request;
35
36pub mod compactor_runner;
37mod context;
38pub mod fast_compactor_runner;
39mod iterator;
40mod shared_buffer_compact;
41pub(super) mod task_progress;
42
43use std::collections::hash_map::Entry;
44use std::collections::{HashMap, VecDeque};
45use std::marker::PhantomData;
46use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
47use std::sync::{Arc, Mutex};
48use std::time::{Duration, SystemTime};
49
50use await_tree::{InstrumentAwait, SpanExt};
51pub use compaction_executor::CompactionExecutor;
52pub use compaction_filter::{
53 CompactionFilter, DummyCompactionFilter, MultiCompactionFilter, TtlCompactionFilter,
54};
55pub use context::{
56 CompactionAwaitTreeRegRef, CompactorContext, await_tree_key, new_compaction_await_tree_reg_ref,
57};
58use futures::{StreamExt, pin_mut};
59use iceberg_compaction::iceberg_compactor_runner::IcebergCompactorRunnerConfigBuilder;
61#[cfg(madsim)]
62pub use iceberg_compaction::set_simulated_pk_index_compaction_result;
63#[cfg(madsim)]
64use iceberg_compaction::simulated_pk_index_compaction_result;
65use iceberg_compaction::{
66 IcebergPlanCompletion, IcebergTaskQueue, IcebergTaskReport, IcebergTaskTracker, PushResult,
67 ReportSendResult, build_drained_iceberg_task_report, build_iceberg_task_report,
68 create_task_execution, flush_pending_iceberg_task_reports, send_or_buffer_iceberg_task_report,
69};
70pub use iterator::{ConcatSstableIterator, SstableStreamIterator};
71use more_asserts::assert_ge;
72use risingwave_hummock_sdk::table_stats::{TableStatsMap, to_prost_table_stats_map};
73use risingwave_hummock_sdk::{
74 HummockCompactionTaskId, HummockSstableObjectId, LocalSstableInfo, compact_task_to_string,
75};
76use risingwave_pb::hummock::compact_task::TaskStatus;
77use risingwave_pb::hummock::subscribe_compaction_event_request::{
78 Event as RequestEvent, HeartBeat, PullTask, ReportTask,
79};
80use risingwave_pb::hummock::subscribe_compaction_event_response::Event as ResponseEvent;
81use risingwave_pb::hummock::{
82 CompactTaskProgress, PbSstableFilterLayout, PbSstableFilterType, ReportCompactionTaskRequest,
83 SubscribeCompactionEventRequest, SubscribeCompactionEventResponse,
84};
85use risingwave_pb::id::IcebergCompactionTaskId;
86use risingwave_rpc_client::HummockMetaClient;
87pub use shared_buffer_compact::compact;
88use tokio::sync::oneshot::Sender;
89use tokio::task::JoinHandle;
90use tokio::time::Instant;
91
92pub use self::compaction_utils::{
93 CompactionStatistics, RemoteBuilderFactory, TaskConfig, check_compaction_result,
94 check_flush_result,
95};
96pub use self::task_progress::TaskProgress;
97use super::multi_builder::CapacitySplitTableBuilder;
98use super::{
99 GetObjectId, HummockError, HummockErrorInner, HummockResult, ObjectIdManager,
100 SstableBuilderOptions, Xor8FilterBuilder, Xor16FilterBuilder,
101};
102use crate::compaction_catalog_manager::{
103 CompactionCatalogAgentRef, CompactionCatalogManager, CompactionCatalogManagerRef,
104};
105use crate::hummock::compactor::compaction_utils::calculate_task_parallelism;
106use crate::hummock::compactor::compactor_runner::{compact_and_build_sst, compact_done};
107use crate::hummock::compactor::iceberg_compaction::TaskKey;
108use crate::hummock::iterator::{Forward, HummockIterator};
109use crate::hummock::{
110 BlockedXor8FilterBuilder, BlockedXor16FilterBuilder, FilterBuilder, NoneFilterBuilder,
111 SharedComapctorObjectIdManager, SstableWriterFactory, StreamingSstableWriterFactory,
112 validate_ssts,
113};
114use crate::monitor::CompactorMetrics;
115
116const COMPACTION_HEARTBEAT_LOG_INTERVAL: Duration = Duration::from_secs(60);
118
119#[derive(Debug, Clone, PartialEq, Eq)]
121struct CompactionLogState {
122 running_parallelism: u32,
123 pull_task_ack: bool,
124 pending_pull_task_count: u32,
125}
126
127#[derive(Debug, Clone, PartialEq, Eq)]
129struct IcebergCompactionLogState {
130 running_parallelism: u32,
131 waiting_parallelism: u32,
132 available_parallelism: u32,
133 pull_task_ack: bool,
134 pending_pull_task_count: u32,
135}
136
137struct LogThrottler<T: PartialEq> {
139 last_logged_state: Option<T>,
140 last_heartbeat: Instant,
141 heartbeat_interval: Duration,
142}
143
144impl<T: PartialEq> LogThrottler<T> {
145 fn new(heartbeat_interval: Duration) -> Self {
146 Self {
147 last_logged_state: None,
148 last_heartbeat: Instant::now(),
149 heartbeat_interval,
150 }
151 }
152
153 fn should_log(&self, current_state: &T) -> bool {
155 self.last_logged_state.as_ref() != Some(current_state)
156 || self.last_heartbeat.elapsed() >= self.heartbeat_interval
157 }
158
159 fn update(&mut self, current_state: T) {
161 self.last_logged_state = Some(current_state);
162 self.last_heartbeat = Instant::now();
163 }
164}
165
166pub struct Compactor {
168 context: CompactorContext,
170 object_id_getter: Arc<dyn GetObjectId>,
171 task_config: TaskConfig,
172 options: SstableBuilderOptions,
173 get_id_time: Arc<AtomicU64>,
174}
175
176pub type CompactOutput = (usize, Vec<LocalSstableInfo>, CompactionStatistics);
177
178impl Compactor {
179 pub fn new(
181 context: CompactorContext,
182 options: SstableBuilderOptions,
183 task_config: TaskConfig,
184 object_id_getter: Arc<dyn GetObjectId>,
185 ) -> Self {
186 Self {
187 context,
188 options,
189 task_config,
190 get_id_time: Arc::new(AtomicU64::new(0)),
191 object_id_getter,
192 }
193 }
194
195 async fn compact_key_range(
200 &self,
201 iter: impl HummockIterator<Direction = Forward>,
202 compaction_filter: impl CompactionFilter,
203 compaction_catalog_agent_ref: CompactionCatalogAgentRef,
204 task_progress: Option<Arc<TaskProgress>>,
205 task_id: Option<HummockCompactionTaskId>,
206 split_index: Option<usize>,
207 ) -> HummockResult<(Vec<LocalSstableInfo>, CompactionStatistics)> {
208 let compact_timer = if self.context.is_share_buffer_compact {
210 self.context
211 .compactor_metrics
212 .write_build_l0_sst_duration
213 .start_timer()
214 } else {
215 self.context
216 .compactor_metrics
217 .compact_sst_duration
218 .start_timer()
219 };
220
221 let (split_table_outputs, table_stats_map) = {
222 let factory = StreamingSstableWriterFactory::new(self.context.sstable_store.clone());
223 match (
224 self.task_config.sstable_filter_type,
225 self.task_config.sstable_filter_layout,
226 ) {
227 (PbSstableFilterType::SstableFilterNone, _) => {
228 self.compact_key_range_impl::<_, NoneFilterBuilder>(
229 factory,
230 iter,
231 compaction_filter,
232 compaction_catalog_agent_ref,
233 task_progress.clone(),
234 self.object_id_getter.clone(),
235 )
236 .instrument_await("compact".verbose())
237 .await?
238 }
239 (PbSstableFilterType::SstableFilterXor8, PbSstableFilterLayout::Blocked) => {
240 self.compact_key_range_impl::<_, BlockedXor8FilterBuilder>(
241 factory,
242 iter,
243 compaction_filter,
244 compaction_catalog_agent_ref,
245 task_progress.clone(),
246 self.object_id_getter.clone(),
247 )
248 .instrument_await("compact".verbose())
249 .await?
250 }
251 (PbSstableFilterType::SstableFilterXor8, PbSstableFilterLayout::Plain) => {
252 self.compact_key_range_impl::<_, Xor8FilterBuilder>(
253 factory,
254 iter,
255 compaction_filter,
256 compaction_catalog_agent_ref,
257 task_progress.clone(),
258 self.object_id_getter.clone(),
259 )
260 .instrument_await("compact".verbose())
261 .await?
262 }
263 (PbSstableFilterType::SstableFilterXor16, PbSstableFilterLayout::Blocked) => {
264 self.compact_key_range_impl::<_, BlockedXor16FilterBuilder>(
265 factory,
266 iter,
267 compaction_filter,
268 compaction_catalog_agent_ref,
269 task_progress.clone(),
270 self.object_id_getter.clone(),
271 )
272 .instrument_await("compact".verbose())
273 .await?
274 }
275 (PbSstableFilterType::SstableFilterXor16, PbSstableFilterLayout::Plain) => {
276 self.compact_key_range_impl::<_, Xor16FilterBuilder>(
277 factory,
278 iter,
279 compaction_filter,
280 compaction_catalog_agent_ref,
281 task_progress.clone(),
282 self.object_id_getter.clone(),
283 )
284 .instrument_await("compact".verbose())
285 .await?
286 }
287 (filter_type, layout) => {
288 return Err(HummockError::compaction_executor(format!(
289 "unresolved SST filter layout in task config: {:?}, {:?}",
290 filter_type, layout
291 )));
292 }
293 }
294 };
295
296 compact_timer.observe_duration();
297
298 Self::report_progress(
299 self.context.compactor_metrics.clone(),
300 task_progress,
301 &split_table_outputs,
302 self.context.is_share_buffer_compact,
303 );
304
305 self.context
306 .compactor_metrics
307 .get_table_id_total_time_duration
308 .observe(self.get_id_time.load(Ordering::Relaxed) as f64 / 1000.0 / 1000.0);
309
310 debug_assert!(
311 split_table_outputs
312 .iter()
313 .all(|table_info| table_info.sst_info.table_ids.is_sorted())
314 );
315
316 if task_id.is_some() {
317 tracing::info!(
319 "Finish Task {:?} split_index {:?} sst count {}",
320 task_id,
321 split_index,
322 split_table_outputs.len()
323 );
324 }
325 Ok((split_table_outputs, table_stats_map))
326 }
327
328 pub fn report_progress(
329 metrics: Arc<CompactorMetrics>,
330 task_progress: Option<Arc<TaskProgress>>,
331 ssts: &Vec<LocalSstableInfo>,
332 is_share_buffer_compact: bool,
333 ) {
334 for sst_info in ssts {
335 let sst_size = sst_info.file_size();
336 if let Some(tracker) = &task_progress {
337 tracker.inc_ssts_uploaded();
338 tracker.dec_num_pending_write_io();
339 }
340 if is_share_buffer_compact {
341 metrics.shared_buffer_to_sstable_size.observe(sst_size as _);
342 } else {
343 metrics.compaction_upload_sst_counts.inc();
344 }
345 }
346 }
347
348 async fn compact_key_range_impl<F: SstableWriterFactory, B: FilterBuilder>(
349 &self,
350 writer_factory: F,
351 iter: impl HummockIterator<Direction = Forward>,
352 compaction_filter: impl CompactionFilter,
353 compaction_catalog_agent_ref: CompactionCatalogAgentRef,
354 task_progress: Option<Arc<TaskProgress>>,
355 object_id_getter: Arc<dyn GetObjectId>,
356 ) -> HummockResult<(Vec<LocalSstableInfo>, CompactionStatistics)> {
357 let builder_factory = RemoteBuilderFactory::<F, B> {
358 object_id_getter,
359 limiter: self.context.memory_limiter.clone(),
360 options: self.options.clone(),
361 policy: self.task_config.cache_policy,
362 remote_rpc_cost: self.get_id_time.clone(),
363 compaction_catalog_agent_ref: compaction_catalog_agent_ref.clone(),
364 sstable_writer_factory: writer_factory,
365 _phantom: PhantomData,
366 };
367
368 let mut sst_builder = CapacitySplitTableBuilder::new(
369 builder_factory,
370 self.context.compactor_metrics.clone(),
371 task_progress.clone(),
372 self.task_config.table_vnode_partition.clone(),
373 self.context
374 .storage_opts
375 .compactor_concurrent_uploading_sst_count,
376 compaction_catalog_agent_ref,
377 );
378 let compaction_statistics = compact_and_build_sst(
379 &mut sst_builder,
380 &self.task_config,
381 self.context.compactor_metrics.clone(),
382 iter,
383 compaction_filter,
384 )
385 .instrument_await("compact_and_build_sst".verbose())
386 .await?;
387
388 let ssts = sst_builder
389 .finish()
390 .instrument_await("builder_finish".verbose())
391 .await?;
392
393 Ok((ssts, compaction_statistics))
394 }
395}
396
397#[must_use]
400pub fn start_iceberg_compactor(
401 compactor_context: CompactorContext,
402 hummock_meta_client: Arc<dyn HummockMetaClient>,
403 iceberg_compaction_memory_limit_bytes: usize,
404) -> (JoinHandle<()>, Sender<()>) {
405 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
406 let stream_retry_interval = Duration::from_secs(30);
407 let periodic_event_update_interval = Duration::from_millis(
408 compactor_context
409 .storage_opts
410 .iceberg_compaction_pull_interval_ms,
411 );
412 let worker_num = compactor_context.compaction_executor.worker_num();
413
414 let max_task_parallelism: u32 = (worker_num as f32
415 * compactor_context.storage_opts.compactor_max_task_multiplier)
416 .ceil() as u32;
417
418 let max_pull_task_count = std::cmp::min(
419 max_task_parallelism,
420 compactor_context
421 .storage_opts
422 .iceberg_compaction_max_pull_task_count,
423 );
424
425 assert_ge!(
426 compactor_context.storage_opts.compactor_max_task_multiplier,
427 0.0
428 );
429
430 let join_handle = tokio::spawn(async move {
431 let pending_parallelism_budget = (max_task_parallelism as f32
433 * compactor_context
434 .storage_opts
435 .iceberg_compaction_pending_parallelism_budget_multiplier)
436 .ceil() as u32;
437 let mut task_queue = IcebergTaskQueue::new_with_metrics(
438 max_task_parallelism,
439 pending_parallelism_budget,
440 iceberg_compaction_memory_limit_bytes,
441 compactor_context.compactor_metrics.clone(),
442 );
443
444 let shutdown_map = Arc::new(Mutex::new(HashMap::<TaskKey, Sender<()>>::new()));
446
447 let (task_completion_tx, mut task_completion_rx) =
449 tokio::sync::mpsc::unbounded_channel::<IcebergPlanCompletion>();
450 let mut task_trackers = HashMap::<IcebergCompactionTaskId, IcebergTaskTracker>::new();
451 let mut pending_task_reports = VecDeque::<IcebergTaskReport>::new();
454
455 let mut min_interval = tokio::time::interval(stream_retry_interval);
456 let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
457
458 let mut log_throttler =
460 LogThrottler::<IcebergCompactionLogState>::new(COMPACTION_HEARTBEAT_LOG_INTERVAL);
461
462 'start_stream: loop {
464 let mut pull_task_ack = true;
467 tokio::select! {
468 _ = min_interval.tick() => {},
470 _ = &mut shutdown_rx => {
472 tracing::info!("Compactor is shutting down");
473 return;
474 }
475 }
476
477 let (request_sender, response_event_stream) = match hummock_meta_client
478 .subscribe_iceberg_compaction_event()
479 .await
480 {
481 Ok((request_sender, response_event_stream)) => {
482 tracing::debug!("Succeeded subscribe_iceberg_compaction_event.");
483 (request_sender, response_event_stream)
484 }
485
486 Err(e) => {
487 tracing::warn!(
488 error = %e.as_report(),
489 "Subscribing to iceberg compaction tasks failed with error. Will retry.",
490 );
491 continue 'start_stream;
492 }
493 };
494
495 if matches!(
496 flush_pending_iceberg_task_reports(&request_sender, &mut pending_task_reports),
497 ReportSendResult::RestartStream
498 ) {
499 continue 'start_stream;
500 }
501
502 pin_mut!(response_event_stream);
503
504 let _executor = compactor_context.compaction_executor.clone();
505
506 let mut event_loop_iteration_now = Instant::now();
508 'consume_stream: loop {
509 {
510 compactor_context
512 .compactor_metrics
513 .compaction_event_loop_iteration_latency
514 .observe(event_loop_iteration_now.elapsed().as_millis() as _);
515 event_loop_iteration_now = Instant::now();
516 }
517
518 let request_sender = request_sender.clone();
519 let event: Option<Result<SubscribeIcebergCompactionEventResponse, _>> = tokio::select! {
520 Some(plan_completion) = task_completion_rx.recv() => {
522 let task_key = plan_completion.task_key;
523 let error_message = plan_completion.error_message;
524 let pk_index_result = plan_completion.pk_index_result;
525 tracing::debug!(
526 task_id = %task_key.0,
527 plan_index = task_key.1,
528 success = error_message.is_none(),
529 "Plan completed, updating queue state"
530 );
531 task_queue.finish_running(task_key);
532
533 let completed_task_id = task_key.0;
534 let Entry::Occupied(mut tracker_entry) =
535 task_trackers.entry(completed_task_id)
536 else {
537 continue 'consume_stream;
538 };
539 tracker_entry
540 .get_mut()
541 .record_completion(error_message, pk_index_result);
542 if !tracker_entry.get().is_finished() {
543 continue 'consume_stream;
544 }
545
546 let tracker = tracker_entry.remove();
547 let sink_id = tracker.sink_id();
548 let admitted_plans = tracker.admitted_plans();
549 let successful_plans = tracker.successful_plans();
550 let failed_plans = tracker.failed_plans();
551 let report = tracker.into_report(completed_task_id);
552 let task_succeeded = report.error_message.is_none();
553 if matches!(
554 send_or_buffer_iceberg_task_report(
555 &request_sender,
556 &mut pending_task_reports,
557 report,
558 ),
559 ReportSendResult::RestartStream
560 ) {
561 continue 'start_stream;
562 }
563 if task_succeeded {
564 tracing::info!(
565 iceberg_component = "compaction_worker",
566 iceberg_operation = "report_task",
567 task_id = %completed_task_id,
568 sink_id = sink_id,
569 admitted_plans = admitted_plans,
570 successful_plans = successful_plans,
571 failed_plans = failed_plans,
572 "iceberg_compaction_task_succeeded",
573 );
574 }
575 continue 'consume_stream;
576 }
577
578 _ = task_queue.wait_schedulable() => {
580 schedule_queued_tasks(
581 &mut task_queue,
582 &compactor_context,
583 &shutdown_map,
584 &task_completion_tx,
585 );
586 continue 'consume_stream;
587 }
588
589 _ = periodic_event_interval.tick() => {
590 let should_restart_stream = handle_meta_task_pulling(
592 &mut pull_task_ack,
593 &task_queue,
594 max_task_parallelism,
595 max_pull_task_count,
596 &request_sender,
597 &mut log_throttler,
598 );
599
600 if should_restart_stream {
601 continue 'start_stream;
602 }
603 continue;
604 }
605 event = response_event_stream.next() => {
606 event
607 }
608
609 _ = &mut shutdown_rx => {
610 tracing::info!("Iceberg Compactor is shutting down");
611 return
612 }
613 };
614
615 match event {
616 Some(Ok(SubscribeIcebergCompactionEventResponse {
617 event,
618 create_at: _create_at,
619 })) => {
620 let event = match event {
621 Some(event) => event,
622 None => continue 'consume_stream,
623 };
624
625 match event {
626 risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::CompactTask(iceberg_compaction_task) => {
627 let task_id = iceberg_compaction_task.task_id;
628 let sink_id = iceberg_compaction_task.sink_id;
629 let is_bounded_round = iceberg_compaction_task.max_file_sequence_number.is_some();
630 #[cfg(madsim)]
635 if iceberg_compaction_task.pk_index_coordinated
636 && fail::eval(
637 "iceberg_pk_index_simulated_compaction_report",
638 |_| (),
639 )
640 .is_some()
641 && let Some(result) = simulated_pk_index_compaction_result()
642 {
643 let mut report = build_iceberg_task_report(task_id, sink_id, None);
644 report.pk_index_result = Some(result);
645 if matches!(
646 send_or_buffer_iceberg_task_report(
647 &request_sender,
648 &mut pending_task_reports,
649 report,
650 ),
651 ReportSendResult::RestartStream
652 ) {
653 continue 'start_stream;
654 }
655 continue 'consume_stream;
656 }
657 let compactor_runner_config = match IcebergCompactorRunnerConfigBuilder::default()
659 .max_parallelism((worker_num as f32 * compactor_context.storage_opts.iceberg_compaction_task_parallelism_ratio) as u32)
660 .min_size_per_partition(compactor_context.storage_opts.iceberg_compaction_min_size_per_partition_mb as u64 * 1024 * 1024)
661 .max_file_count_per_partition(compactor_context.storage_opts.iceberg_compaction_max_file_count_per_partition)
662 .enable_validate_compaction(compactor_context.storage_opts.iceberg_compaction_enable_validate)
663 .max_record_batch_rows(compactor_context.storage_opts.iceberg_compaction_max_record_batch_rows)
664 .enable_heuristic_output_parallelism(compactor_context.storage_opts.iceberg_compaction_enable_heuristic_output_parallelism)
665 .max_concurrent_closes(compactor_context.storage_opts.iceberg_compaction_max_concurrent_closes)
666 .enable_prefetch(compactor_context.storage_opts.iceberg_compaction_enable_prefetch)
667 .target_binpack_group_size_mb(
668 compactor_context.storage_opts.iceberg_compaction_target_binpack_group_size_mb
669 )
670 .min_group_size_mb(
671 compactor_context.storage_opts.iceberg_compaction_min_group_size_mb
672 )
673 .min_group_file_count(
674 compactor_context.storage_opts.iceberg_compaction_min_group_file_count
675 )
676 .build() {
677 Ok(config) => config,
678 Err(e) => {
679 tracing::warn!(
680 iceberg_component = "compaction_worker",
681 iceberg_operation = "build_runner_config",
682 error = %e.as_report(),
683 task_id = %task_id,
684 sink_id = sink_id,
685 "iceberg_compaction_runner_config_failed",
686 );
687 let report = build_iceberg_task_report(
688 task_id,
689 sink_id,
690 Some(format!(
691 "Failed to build iceberg compactor runner config: {}",
692 e.as_report()
693 )),
694 );
695 if matches!(
696 send_or_buffer_iceberg_task_report(
697 &request_sender,
698 &mut pending_task_reports,
699 report,
700 ),
701 ReportSendResult::RestartStream
702 ) {
703 continue 'start_stream;
704 }
705 continue 'consume_stream;
706 }
707 };
708
709 let task_execution = match create_task_execution(
711 iceberg_compaction_task,
712 compactor_runner_config,
713 compactor_context.compactor_metrics.clone(),
714 ).await {
715 Ok(task_execution) => task_execution,
716 Err(e) => {
717 tracing::warn!(
718 iceberg_component = "compaction_worker",
719 iceberg_operation = "plan_task",
720 error = %e.as_report(),
721 task_id = %task_id,
722 sink_id = sink_id,
723 "iceberg_compaction_task_plan_failed",
724 );
725 let report = build_iceberg_task_report(
726 task_id,
727 sink_id,
728 Some(format!(
729 "Failed to create iceberg compaction task execution: {}",
730 e.as_report()
731 )),
732 );
733 if matches!(
734 send_or_buffer_iceberg_task_report(
735 &request_sender,
736 &mut pending_task_reports,
737 report,
738 ),
739 ReportSendResult::RestartStream
740 ) {
741 continue 'start_stream;
742 }
743 continue 'consume_stream;
744 }
745 };
746
747 let sink_id = task_execution.sink_id;
748 let plan_runners = task_execution.plan_runners;
749
750 if plan_runners.is_empty() {
751 tracing::info!(
753 iceberg_component = "compaction_worker",
754 iceberg_operation = "enqueue_plan",
755 task_id = %task_id,
756 sink_id = sink_id,
757 "iceberg_compaction_task_skipped_no_plans",
758 );
759 let report = if is_bounded_round {
760 build_drained_iceberg_task_report(task_id, sink_id)
761 } else {
762 build_iceberg_task_report(task_id, sink_id, None)
763 };
764 if matches!(
765 send_or_buffer_iceberg_task_report(
766 &request_sender,
767 &mut pending_task_reports,
768 report,
769 ),
770 ReportSendResult::RestartStream
771 ) {
772 continue 'start_stream;
773 }
774 continue 'consume_stream;
775 }
776
777 let planned_count = plan_runners.len();
779 let mut admitted_count = 0;
780 let mut rejection_reason = None;
781
782 for runner in plan_runners {
783 let meta = runner.to_meta();
784 let plan_index = meta.plan_index;
785 let required_parallelism = runner.required_parallelism();
786 let memory_reservation_bytes = meta.memory_reservation_bytes;
787 let runner_sink_id = runner.sink_id;
788 let runner_compaction_kind = runner.compaction_kind();
789 let runner_table = runner.table_ident.to_string();
790 let push_result = task_queue.push(meta, Some(runner));
791
792 match push_result {
793 PushResult::Added => {
794 admitted_count += 1;
795 tracing::debug!(
796 iceberg_component = "compaction_worker",
797 iceberg_operation = "enqueue_plan",
798 task_id = %task_id,
799 sink_id = runner_sink_id,
800 plan_index = plan_index,
801 task_type = ?runner_compaction_kind,
802 table = %runner_table,
803 required_parallelism = required_parallelism,
804 "iceberg_compaction_plan_enqueued",
805 );
806 },
807 PushResult::RejectedCapacity => {
808 tracing::warn!(
809 iceberg_component = "compaction_worker",
810 iceberg_operation = "enqueue_plan",
811 task_id = %task_id,
812 sink_id = runner_sink_id,
813 plan_index = plan_index,
814 task_type = ?runner_compaction_kind,
815 table = %runner_table,
816 required_parallelism = required_parallelism,
817 pending_budget = pending_parallelism_budget,
818 admitted_count = admitted_count,
819 planned_count = planned_count,
820 "iceberg_compaction_plan_rejected_capacity",
821 );
822 break;
824 },
825 PushResult::RejectedTooLarge => {
826 rejection_reason.get_or_insert_with(|| {
827 format!(
828 "Iceberg compaction plan {plan_index} requires parallelism {required_parallelism} and an estimated {memory_reservation_bytes} bytes of memory, exceeding worker capacities of {max_task_parallelism} parallelism and {iceberg_compaction_memory_limit_bytes} bytes"
829 )
830 });
831 tracing::error!(
832 iceberg_component = "compaction_worker",
833 iceberg_operation = "enqueue_plan",
834 task_id = %task_id,
835 sink_id = runner_sink_id,
836 plan_index = plan_index,
837 task_type = ?runner_compaction_kind,
838 table = %runner_table,
839 required_parallelism = required_parallelism,
840 max_parallelism = max_task_parallelism,
841 memory_reservation_bytes,
842 memory_budget_bytes = iceberg_compaction_memory_limit_bytes,
843 "iceberg_compaction_plan_rejected_too_large",
844 );
845 },
846 PushResult::RejectedInvalidParallelism => {
847 tracing::error!(
848 iceberg_component = "compaction_worker",
849 iceberg_operation = "enqueue_plan",
850 task_id = %task_id,
851 sink_id = runner_sink_id,
852 plan_index = plan_index,
853 task_type = ?runner_compaction_kind,
854 table = %runner_table,
855 required_parallelism = required_parallelism,
856 "iceberg_compaction_plan_rejected_invalid_parallelism",
857 );
858 },
859 PushResult::RejectedDuplicate => {
860 tracing::error!(
861 iceberg_component = "compaction_worker",
862 iceberg_operation = "enqueue_plan",
863 task_id = %task_id,
864 sink_id = runner_sink_id,
865 plan_index = plan_index,
866 task_type = ?runner_compaction_kind,
867 table = %runner_table,
868 "iceberg_compaction_plan_rejected_duplicate",
869 );
870 }
871 }
872 }
873
874 if admitted_count == 0 {
875 let report = build_iceberg_task_report(
876 task_id,
877 sink_id,
878 Some(rejection_reason.unwrap_or_else(|| {
879 "Failed to enqueue all iceberg compaction plans".to_owned()
880 })),
881 );
882 if matches!(
883 send_or_buffer_iceberg_task_report(
884 &request_sender,
885 &mut pending_task_reports,
886 report,
887 ),
888 ReportSendResult::RestartStream
889 ) {
890 continue 'start_stream;
891 }
892 } else {
893 let fully_admitted_bounded_round =
897 is_bounded_round && admitted_count == planned_count;
898 task_trackers.insert(
899 task_id,
900 IcebergTaskTracker::new(
901 sink_id,
902 admitted_count,
903 fully_admitted_bounded_round,
904 ),
905 );
906 }
907
908 tracing::info!(
909 iceberg_component = "compaction_worker",
910 iceberg_operation = "enqueue_plan",
911 task_id = %task_id,
912 sink_id = sink_id,
913 planned_count = planned_count,
914 admitted_count = admitted_count,
915 "iceberg_compaction_task_enqueue_finished",
916 );
917 },
918 risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::PullTaskAck(_) => {
919 pull_task_ack = true;
921 },
922 risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::CancelCompactTask(cancel_compact_task) => {
923 cancel_iceberg_task(
924 cancel_compact_task.task_id,
925 &mut task_queue,
926 &shutdown_map,
927 &mut task_trackers,
928 );
929 },
930 }
931 }
932 Some(Err(e)) => {
933 tracing::warn!("Failed to consume stream. {}", e.message());
934 continue 'start_stream;
935 }
936 _ => {
937 continue 'start_stream;
939 }
940 }
941 }
942 }
943 });
944
945 (join_handle, shutdown_tx)
946}
947
948#[must_use]
951pub fn start_compactor(
952 compactor_context: CompactorContext,
953 hummock_meta_client: Arc<dyn HummockMetaClient>,
954 object_id_manager: Arc<ObjectIdManager>,
955 compaction_catalog_manager_ref: CompactionCatalogManagerRef,
956) -> (JoinHandle<()>, Sender<()>) {
957 type CompactionShutdownMap = Arc<Mutex<HashMap<u64, Sender<()>>>>;
958 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
959 let stream_retry_interval = Duration::from_secs(30);
960 let task_progress = compactor_context.task_progress_manager.clone();
961 let periodic_event_update_interval = Duration::from_millis(1000);
962
963 let max_task_parallelism: u32 = (compactor_context.compaction_executor.worker_num() as f32
964 * compactor_context.storage_opts.compactor_max_task_multiplier)
965 .ceil() as u32;
966 let running_task_parallelism = Arc::new(AtomicU32::new(0));
967
968 const MAX_PULL_TASK_COUNT: u32 = 4;
969 let max_pull_task_count = std::cmp::min(max_task_parallelism, MAX_PULL_TASK_COUNT);
970
971 assert_ge!(
972 compactor_context.storage_opts.compactor_max_task_multiplier,
973 0.0
974 );
975
976 let join_handle = tokio::spawn(async move {
977 let shutdown_map = CompactionShutdownMap::default();
978 let mut min_interval = tokio::time::interval(stream_retry_interval);
979 let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
980
981 let mut log_throttler =
983 LogThrottler::<CompactionLogState>::new(COMPACTION_HEARTBEAT_LOG_INTERVAL);
984
985 'start_stream: loop {
987 let mut pull_task_ack = true;
990 tokio::select! {
991 _ = min_interval.tick() => {},
993 _ = &mut shutdown_rx => {
995 tracing::info!("Compactor is shutting down");
996 return;
997 }
998 }
999
1000 let (request_sender, response_event_stream) =
1001 match hummock_meta_client.subscribe_compaction_event().await {
1002 Ok((request_sender, response_event_stream)) => {
1003 tracing::debug!("Succeeded subscribe_compaction_event.");
1004 (request_sender, response_event_stream)
1005 }
1006
1007 Err(e) => {
1008 tracing::warn!(
1009 error = %e.as_report(),
1010 "Subscribing to compaction tasks failed with error. Will retry.",
1011 );
1012 continue 'start_stream;
1013 }
1014 };
1015
1016 pin_mut!(response_event_stream);
1017
1018 let executor = compactor_context.compaction_executor.clone();
1019 let object_id_manager = object_id_manager.clone();
1020
1021 let mut event_loop_iteration_now = Instant::now();
1023 'consume_stream: loop {
1024 {
1025 compactor_context
1027 .compactor_metrics
1028 .compaction_event_loop_iteration_latency
1029 .observe(event_loop_iteration_now.elapsed().as_millis() as _);
1030 event_loop_iteration_now = Instant::now();
1031 }
1032
1033 let running_task_parallelism = running_task_parallelism.clone();
1034 let request_sender = request_sender.clone();
1035 let event: Option<Result<SubscribeCompactionEventResponse, _>> = tokio::select! {
1036 _ = periodic_event_interval.tick() => {
1037 let progress_list = get_task_progress(task_progress.clone());
1038
1039 if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1040 event: Some(RequestEvent::HeartBeat(
1041 HeartBeat {
1042 progress: progress_list
1043 }
1044 )),
1045 create_at: SystemTime::now()
1046 .duration_since(std::time::UNIX_EPOCH)
1047 .expect("Clock may have gone backwards")
1048 .as_millis() as u64,
1049 }) {
1050 tracing::warn!(error = %e.as_report(), "Failed to report task progress");
1051 continue 'start_stream;
1053 }
1054
1055
1056 let mut pending_pull_task_count = 0;
1057 if pull_task_ack {
1058 pending_pull_task_count = (max_task_parallelism - running_task_parallelism.load(Ordering::SeqCst)).min(max_pull_task_count);
1060
1061 if pending_pull_task_count > 0 {
1062 if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1063 event: Some(RequestEvent::PullTask(
1064 PullTask {
1065 pull_task_count: pending_pull_task_count,
1066 }
1067 )),
1068 create_at: SystemTime::now()
1069 .duration_since(std::time::UNIX_EPOCH)
1070 .expect("Clock may have gone backwards")
1071 .as_millis() as u64,
1072 }) {
1073 tracing::warn!(error = %e.as_report(), "Failed to pull task");
1074
1075 continue 'start_stream;
1077 } else {
1078 pull_task_ack = false;
1079 }
1080 }
1081 }
1082
1083 let running_count = running_task_parallelism.load(Ordering::SeqCst);
1084 let current_state = CompactionLogState {
1085 running_parallelism: running_count,
1086 pull_task_ack,
1087 pending_pull_task_count,
1088 };
1089
1090 if log_throttler.should_log(¤t_state) {
1092 tracing::info!(
1093 running_parallelism_count = %current_state.running_parallelism,
1094 pull_task_ack = %current_state.pull_task_ack,
1095 pending_pull_task_count = %current_state.pending_pull_task_count
1096 );
1097 log_throttler.update(current_state);
1098 }
1099
1100 continue;
1101 }
1102 event = response_event_stream.next() => {
1103 event
1104 }
1105
1106 _ = &mut shutdown_rx => {
1107 tracing::info!("Compactor is shutting down");
1108 return
1109 }
1110 };
1111
1112 fn send_report_task_event(
1113 compact_task: &CompactTask,
1114 table_stats: TableStatsMap,
1115 object_timestamps: HashMap<HummockSstableObjectId, u64>,
1116 request_sender: &mpsc::UnboundedSender<SubscribeCompactionEventRequest>,
1117 ) {
1118 if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1119 event: Some(RequestEvent::ReportTask(ReportTask {
1120 task_id: compact_task.task_id,
1121 task_status: compact_task.task_status.into(),
1122 sorted_output_ssts: compact_task
1123 .sorted_output_ssts
1124 .iter()
1125 .map(|sst| sst.into())
1126 .collect(),
1127 table_stats_change: to_prost_table_stats_map(table_stats),
1128 object_timestamps,
1129 })),
1130 create_at: SystemTime::now()
1131 .duration_since(std::time::UNIX_EPOCH)
1132 .expect("Clock may have gone backwards")
1133 .as_millis() as u64,
1134 }) {
1135 let task_id = compact_task.task_id;
1136 tracing::warn!(error = %e.as_report(), "Failed to report task {task_id:?}");
1137 }
1138 }
1139
1140 match event {
1141 Some(Ok(SubscribeCompactionEventResponse { event, create_at })) => {
1142 let event = match event {
1143 Some(event) => event,
1144 None => continue 'consume_stream,
1145 };
1146 let shutdown = shutdown_map.clone();
1147 let context = compactor_context.clone();
1148 let consumed_latency_ms = SystemTime::now()
1149 .duration_since(std::time::UNIX_EPOCH)
1150 .expect("Clock may have gone backwards")
1151 .as_millis() as u64
1152 - create_at;
1153 context
1154 .compactor_metrics
1155 .compaction_event_consumed_latency
1156 .observe(consumed_latency_ms as _);
1157
1158 let object_id_manager = object_id_manager.clone();
1159 let compaction_catalog_manager_ref = compaction_catalog_manager_ref.clone();
1160
1161 match event {
1162 ResponseEvent::CompactTask(compact_task) => {
1163 let compact_task = CompactTask::from(compact_task);
1164 let parallelism =
1165 calculate_task_parallelism(&compact_task, &context);
1166
1167 assert_ne!(parallelism, 0, "splits cannot be empty");
1168
1169 if (max_task_parallelism
1170 - running_task_parallelism.load(Ordering::SeqCst))
1171 < parallelism as u32
1172 {
1173 tracing::warn!(
1174 "Not enough core parallelism to serve the task {} task_parallelism {} running_task_parallelism {} max_task_parallelism {}",
1175 compact_task.task_id,
1176 parallelism,
1177 max_task_parallelism,
1178 running_task_parallelism.load(Ordering::Relaxed),
1179 );
1180 let (compact_task, table_stats, object_timestamps) =
1181 compact_done(
1182 compact_task,
1183 context.clone(),
1184 vec![],
1185 TaskStatus::NoAvailCpuResourceCanceled,
1186 );
1187
1188 send_report_task_event(
1189 &compact_task,
1190 table_stats,
1191 object_timestamps,
1192 &request_sender,
1193 );
1194
1195 continue 'consume_stream;
1196 }
1197
1198 running_task_parallelism
1199 .fetch_add(parallelism as u32, Ordering::SeqCst);
1200 executor.spawn(async move {
1201 let (tx, rx) = tokio::sync::oneshot::channel();
1202 let task_id = compact_task.task_id;
1203 shutdown.lock().unwrap().insert(task_id, tx);
1204
1205 let (
1206 (compact_task, table_stats, object_timestamps),
1207 _memory_tracker,
1208 ) = compactor_runner::compact(
1209 context.clone(),
1210 compact_task,
1211 rx,
1212 object_id_manager.clone(),
1213 compaction_catalog_manager_ref.clone(),
1214 )
1215 .await;
1216
1217 shutdown.lock().unwrap().remove(&task_id);
1218 running_task_parallelism
1219 .fetch_sub(parallelism as u32, Ordering::SeqCst);
1220
1221 send_report_task_event(
1222 &compact_task,
1223 table_stats,
1224 object_timestamps,
1225 &request_sender,
1226 );
1227
1228 let enable_check_compaction_result =
1229 context.storage_opts.check_compaction_result;
1230 let need_check_task =
1231 !compact_task.sorted_output_ssts.is_empty()
1232 && compact_task.task_status == TaskStatus::Success;
1233
1234 if enable_check_compaction_result && need_check_task {
1235 let read_table_ids = compact_task
1236 .get_table_ids_from_input_ssts()
1237 .collect::<Vec<_>>();
1238 match compaction_catalog_manager_ref.acquire(read_table_ids).await {
1239 Ok(compaction_catalog_agent_ref) => {
1240 match check_compaction_result(&compact_task, context.clone(), compaction_catalog_agent_ref).await
1241 {
1242 Err(e) => {
1243 tracing::warn!(error = %e.as_report(), "Failed to check compaction task {}",compact_task.task_id);
1244 }
1245 Ok(true) => (),
1246 Ok(false) => {
1247 panic!("Failed to pass consistency check for result of compaction task:\n{:?}", compact_task_to_string(&compact_task));
1248 }
1249 }
1250 },
1251 Err(e) => {
1252 tracing::warn!(error = %e.as_report(), "failed to acquire compaction catalog agent");
1253 }
1254 }
1255 }
1256 });
1257 }
1258 #[expect(deprecated)]
1259 ResponseEvent::VacuumTask(_) => {
1260 unreachable!("unexpected vacuum task");
1261 }
1262 #[expect(deprecated)]
1263 ResponseEvent::FullScanTask(_) => {
1264 unreachable!("unexpected scan task");
1265 }
1266 #[expect(deprecated)]
1267 ResponseEvent::ValidationTask(validation_task) => {
1268 let validation_task = ValidationTask::from(validation_task);
1269 executor.spawn(async move {
1270 validate_ssts(validation_task, context.sstable_store.clone())
1271 .await;
1272 });
1273 }
1274 ResponseEvent::CancelCompactTask(cancel_compact_task) => match shutdown
1275 .lock()
1276 .unwrap()
1277 .remove(&cancel_compact_task.task_id)
1278 {
1279 Some(tx) => {
1280 if tx.send(()).is_err() {
1281 tracing::warn!(
1282 "Cancellation of compaction task failed. task_id: {}",
1283 cancel_compact_task.task_id
1284 );
1285 }
1286 }
1287 _ => {
1288 tracing::warn!(
1289 "Attempting to cancel non-existent compaction task. task_id: {}",
1290 cancel_compact_task.task_id
1291 );
1292 }
1293 },
1294
1295 ResponseEvent::PullTaskAck(_pull_task_ack) => {
1296 pull_task_ack = true;
1298 }
1299 }
1300 }
1301 Some(Err(e)) => {
1302 tracing::warn!("Failed to consume stream. {}", e.message());
1303 continue 'start_stream;
1304 }
1305 _ => {
1306 continue 'start_stream;
1308 }
1309 }
1310 }
1311 }
1312 });
1313
1314 (join_handle, shutdown_tx)
1315}
1316
1317#[must_use]
1320pub fn start_shared_compactor(
1321 grpc_proxy_client: GrpcCompactorProxyClient,
1322 mut receiver: mpsc::UnboundedReceiver<Request<DispatchCompactionTaskRequest>>,
1323 context: CompactorContext,
1324) -> (JoinHandle<()>, Sender<()>) {
1325 type CompactionShutdownMap = Arc<Mutex<HashMap<u64, Sender<()>>>>;
1326 let task_progress = context.task_progress_manager.clone();
1327 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
1328 let periodic_event_update_interval = Duration::from_millis(1000);
1329
1330 let join_handle = tokio::spawn(async move {
1331 let shutdown_map = CompactionShutdownMap::default();
1332
1333 let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
1334 let executor = context.compaction_executor.clone();
1335 let report_heartbeat_client = grpc_proxy_client.clone();
1336 'consume_stream: loop {
1337 let request: Option<Request<DispatchCompactionTaskRequest>> = tokio::select! {
1338 _ = periodic_event_interval.tick() => {
1339 let progress_list = get_task_progress(task_progress.clone());
1340 let report_compaction_task_request = ReportCompactionTaskRequest{
1341 event: Some(ReportCompactionTaskEvent::HeartBeat(
1342 SharedHeartBeat {
1343 progress: progress_list
1344 }
1345 )),
1346 };
1347 if let Err(e) = report_heartbeat_client.report_compaction_task(report_compaction_task_request).await{
1348 tracing::warn!(error = %e.as_report(), "Failed to report heartbeat");
1349 }
1350 continue
1351 }
1352
1353
1354 _ = &mut shutdown_rx => {
1355 tracing::info!("Compactor is shutting down");
1356 return
1357 }
1358
1359 request = receiver.recv() => {
1360 request
1361 }
1362
1363 };
1364 match request {
1365 Some(request) => {
1366 let context = context.clone();
1367 let shutdown = shutdown_map.clone();
1368
1369 let cloned_grpc_proxy_client = grpc_proxy_client.clone();
1370 executor.spawn(async move {
1371 let DispatchCompactionTaskRequest {
1372 tables,
1373 output_object_ids,
1374 task: dispatch_task,
1375 } = request.into_inner();
1376 let table_id_to_catalog = tables.into_iter().fold(HashMap::new(), |mut acc, table| {
1377 acc.insert(table.id, table);
1378 acc
1379 });
1380
1381 let mut output_object_ids_deque: VecDeque<_> = VecDeque::new();
1382 output_object_ids_deque.extend(output_object_ids.into_iter().map(Into::<HummockSstableObjectId>::into));
1383 let shared_compactor_object_id_manager =
1384 SharedComapctorObjectIdManager::new(output_object_ids_deque, cloned_grpc_proxy_client.clone(), context.storage_opts.sstable_id_remote_fetch_number);
1385 match dispatch_task.unwrap() {
1386 dispatch_compaction_task_request::Task::CompactTask(compact_task) => {
1387 let compact_task = CompactTask::from(&compact_task);
1388 let (tx, rx) = tokio::sync::oneshot::channel();
1389 let task_id = compact_task.task_id;
1390 shutdown.lock().unwrap().insert(task_id, tx);
1391
1392 let compaction_catalog_manager_ref =
1393 Arc::new(CompactionCatalogManager::new_preloaded(
1394 table_id_to_catalog,
1395 ));
1396 let ((compact_task, table_stats, object_timestamps), _memory_tracker)= compactor_runner::compact(
1397 context.clone(),
1398 compact_task,
1399 rx,
1400 shared_compactor_object_id_manager,
1401 compaction_catalog_manager_ref.clone(),
1402 )
1403 .await;
1404 shutdown.lock().unwrap().remove(&task_id);
1405 let report_compaction_task_request = ReportCompactionTaskRequest {
1406 event: Some(ReportCompactionTaskEvent::ReportTask(ReportSharedTask {
1407 compact_task: Some(PbCompactTask::from(&compact_task)),
1408 table_stats_change: to_prost_table_stats_map(table_stats),
1409 object_timestamps,
1410 })),
1411 };
1412
1413 match cloned_grpc_proxy_client
1414 .report_compaction_task(report_compaction_task_request)
1415 .await
1416 {
1417 Ok(_) => {
1418 let enable_check_compaction_result =
1420 context.storage_opts.check_compaction_result;
1421 let need_check_task =
1422 !compact_task.sorted_output_ssts.is_empty()
1423 && compact_task.task_status
1424 == TaskStatus::Success;
1425 if enable_check_compaction_result && need_check_task {
1426 let read_table_ids = compact_task
1427 .get_table_ids_from_input_ssts()
1428 .collect::<Vec<_>>();
1429 match compaction_catalog_manager_ref.acquire(read_table_ids).await {
1430 Ok(compaction_catalog_agent_ref) => {
1431 match check_compaction_result(&compact_task, context.clone(), compaction_catalog_agent_ref).await
1432 {
1433 Err(e) => {
1434 tracing::warn!(error = %e.as_report(), "Failed to check compaction task {}", task_id);
1435 }
1436 Ok(true) => (),
1437 Ok(false) => {
1438 panic!("Failed to pass consistency check for result of compaction task:\n{:?}", compact_task_to_string(&compact_task));
1439 }
1440 }
1441 }
1442 Err(e) => {
1443 tracing::warn!(error = %e.as_report(), "failed to acquire compaction catalog agent");
1444 }
1445 }
1446 }
1447 }
1448 Err(e) => tracing::warn!(error = %e.as_report(), "Failed to report task {task_id:?}"),
1449 }
1450
1451 }
1452 dispatch_compaction_task_request::Task::VacuumTask(_) => {
1453 unreachable!("unexpected vacuum task");
1454 }
1455 dispatch_compaction_task_request::Task::FullScanTask(_) => {
1456 unreachable!("unexpected scan task");
1457 }
1458 dispatch_compaction_task_request::Task::ValidationTask(validation_task) => {
1459 let validation_task = ValidationTask::from(validation_task);
1460 validate_ssts(validation_task, context.sstable_store.clone()).await;
1461 }
1462 dispatch_compaction_task_request::Task::CancelCompactTask(cancel_compact_task) => {
1463 match shutdown
1464 .lock()
1465 .unwrap()
1466 .remove(&cancel_compact_task.task_id)
1467 { Some(tx) => {
1468 if tx.send(()).is_err() {
1469 tracing::warn!(
1470 "Cancellation of compaction task failed. task_id: {}",
1471 cancel_compact_task.task_id
1472 );
1473 }
1474 } _ => {
1475 tracing::warn!(
1476 "Attempting to cancel non-existent compaction task. task_id: {}",
1477 cancel_compact_task.task_id
1478 );
1479 }}
1480 }
1481 }
1482 });
1483 }
1484 None => continue 'consume_stream,
1485 }
1486 }
1487 });
1488 (join_handle, shutdown_tx)
1489}
1490
1491fn get_task_progress(
1492 task_progress: Arc<
1493 parking_lot::lock_api::Mutex<parking_lot::RawMutex, HashMap<u64, Arc<TaskProgress>>>,
1494 >,
1495) -> Vec<CompactTaskProgress> {
1496 let mut progress_list = Vec::new();
1497 for (&task_id, progress) in &*task_progress.lock() {
1498 progress_list.push(progress.snapshot(task_id));
1499 }
1500 progress_list
1501}
1502
1503fn schedule_queued_tasks(
1505 task_queue: &mut IcebergTaskQueue,
1506 compactor_context: &CompactorContext,
1507 shutdown_map: &Arc<Mutex<HashMap<TaskKey, Sender<()>>>>,
1508 task_completion_tx: &tokio::sync::mpsc::UnboundedSender<IcebergPlanCompletion>,
1509) {
1510 while let Some(popped_task) = task_queue.pop() {
1511 let task_id = popped_task.meta.task_id;
1512 let plan_index = popped_task.meta.plan_index;
1513 let task_key = (task_id, plan_index);
1514
1515 let unique_ident = popped_task.runner.as_ref().map(|r| r.unique_ident());
1517
1518 let Some(runner) = popped_task.runner else {
1519 tracing::error!(
1520 iceberg_component = "compaction_worker",
1521 iceberg_operation = "schedule_plan",
1522 task_id = %task_id,
1523 plan_index = plan_index,
1524 "iceberg_compaction_plan_missing_runner",
1525 );
1526 task_queue.finish_running(task_key);
1527 continue;
1528 };
1529 let runner_sink_id = runner.sink_id;
1530 let runner_compaction_kind = runner.compaction_kind();
1531 let runner_table = runner.table_ident.to_string();
1532
1533 let executor = compactor_context.compaction_executor.clone();
1534 let shutdown_map_clone = shutdown_map.clone();
1535 let completion_tx_clone = task_completion_tx.clone();
1536 let (tx, rx) = tokio::sync::oneshot::channel();
1537
1538 {
1539 let mut shutdown_guard = shutdown_map.lock().unwrap();
1540 shutdown_guard.insert(task_key, tx);
1541 }
1542
1543 tracing::info!(
1544 iceberg_component = "compaction_worker",
1545 iceberg_operation = "schedule_plan",
1546 task_id = %task_id,
1547 sink_id = runner_sink_id,
1548 plan_index = plan_index,
1549 task_type = ?runner_compaction_kind,
1550 table = %runner_table,
1551 unique_ident = ?unique_ident,
1552 required_parallelism = popped_task.meta.required_parallelism,
1553 "iceberg_compaction_plan_started_from_queue",
1554 );
1555
1556 executor.spawn(async move {
1557 let _cleanup_guard = scopeguard::guard(shutdown_map_clone, move |shutdown_map| {
1558 let mut shutdown_guard = shutdown_map.lock().unwrap();
1559 shutdown_guard.remove(&task_key);
1560 });
1561
1562 let result = Box::pin(runner.compact(rx)).await;
1563
1564 let completion = match result {
1565 Ok((_stats, pk_index_result)) => IcebergPlanCompletion {
1566 task_key,
1567 error_message: None,
1568 pk_index_result,
1569 },
1570 Err(e) => {
1571 if is_cancelled_iceberg_compaction_error(&e) {
1572 tracing::info!(
1573 iceberg_component = "compaction_worker",
1574 iceberg_operation = "execute_plan",
1575 task_id = %task_key.0,
1576 sink_id = runner_sink_id,
1577 plan_index = task_key.1,
1578 task_type = ?runner_compaction_kind,
1579 table = %runner_table,
1580 "iceberg_compaction_plan_cancelled",
1581 );
1582 } else {
1583 tracing::warn!(
1584 iceberg_component = "compaction_worker",
1585 iceberg_operation = "execute_plan",
1586 error = %e.as_report(),
1587 task_id = %task_key.0,
1588 sink_id = runner_sink_id,
1589 plan_index = task_key.1,
1590 task_type = ?runner_compaction_kind,
1591 table = %runner_table,
1592 "iceberg_compaction_plan_failed",
1593 );
1594 }
1595 IcebergPlanCompletion {
1596 task_key,
1597 error_message: Some(e.to_report_string()),
1598 pk_index_result: None,
1599 }
1600 }
1601 };
1602
1603 if completion_tx_clone.send(completion).is_err() {
1604 tracing::warn!(
1605 iceberg_component = "compaction_worker",
1606 iceberg_operation = "notify_plan_completion",
1607 task_id = %task_key.0,
1608 sink_id = runner_sink_id,
1609 plan_index = task_key.1,
1610 task_type = ?runner_compaction_kind,
1611 table = %runner_table,
1612 "iceberg_compaction_plan_completion_send_failed",
1613 );
1614 }
1615 });
1616 }
1617}
1618
1619fn is_cancelled_iceberg_compaction_error(error: &crate::hummock::HummockError) -> bool {
1620 matches!(
1621 error.inner(),
1622 HummockErrorInner::CompactionExecutor(message) if message == "Plan cancelled"
1623 )
1624}
1625
1626fn cancel_iceberg_task(
1627 task_id: IcebergCompactionTaskId,
1628 task_queue: &mut IcebergTaskQueue,
1629 shutdown_map: &Arc<Mutex<HashMap<TaskKey, Sender<()>>>>,
1630 task_trackers: &mut HashMap<IcebergCompactionTaskId, IcebergTaskTracker>,
1631) {
1632 let cancelled_waiting = task_queue.cancel_waiting_task(task_id);
1637 let removed_tracker = task_trackers.remove(&task_id).is_some();
1638
1639 let cancelled_running = {
1640 let mut shutdown_guard = shutdown_map.lock().unwrap();
1641 let task_keys: Vec<_> = shutdown_guard
1642 .keys()
1643 .filter(|(running_task_id, _)| *running_task_id == task_id)
1644 .copied()
1645 .collect();
1646
1647 for task_key in &task_keys {
1648 if let Some(tx) = shutdown_guard.remove(task_key)
1649 && tx.send(()).is_err()
1650 {
1651 tracing::debug!(
1652 task_id = %task_key.0,
1653 plan_index = task_key.1,
1654 "Iceberg compaction plan shutdown receiver already closed during cancellation"
1655 );
1656 }
1657 }
1658
1659 task_keys.len()
1660 };
1661
1662 if cancelled_waiting == 0 && cancelled_running == 0 && !removed_tracker {
1663 tracing::warn!(
1664 task_id = %task_id,
1665 "Attempting to cancel non-existent iceberg compaction task"
1666 );
1667 } else {
1668 tracing::info!(
1669 task_id = %task_id,
1670 cancelled_waiting = cancelled_waiting,
1671 cancelled_running = cancelled_running,
1672 removed_tracker = removed_tracker,
1673 "Cancelled iceberg compaction task"
1674 );
1675 }
1676}
1677
1678fn handle_meta_task_pulling(
1681 pull_task_ack: &mut bool,
1682 task_queue: &IcebergTaskQueue,
1683 max_task_parallelism: u32,
1684 max_pull_task_count: u32,
1685 request_sender: &mpsc::UnboundedSender<SubscribeIcebergCompactionEventRequest>,
1686 log_throttler: &mut LogThrottler<IcebergCompactionLogState>,
1687) -> bool {
1688 let mut pending_pull_task_count = 0;
1689 if *pull_task_ack {
1690 let current_reserved_parallelism = task_queue
1695 .running_parallelism_sum()
1696 .saturating_add(task_queue.waiting_parallelism_sum());
1697 pending_pull_task_count = max_task_parallelism
1698 .saturating_sub(current_reserved_parallelism)
1699 .min(max_pull_task_count);
1700
1701 if pending_pull_task_count > 0 {
1702 if let Err(e) = request_sender.send(SubscribeIcebergCompactionEventRequest {
1703 event: Some(subscribe_iceberg_compaction_event_request::Event::PullTask(
1704 subscribe_iceberg_compaction_event_request::PullTask {
1705 pull_task_count: pending_pull_task_count,
1706 },
1707 )),
1708 create_at: SystemTime::now()
1709 .duration_since(std::time::UNIX_EPOCH)
1710 .expect("Clock may have gone backwards")
1711 .as_millis() as u64,
1712 }) {
1713 tracing::warn!(error = %e.as_report(), "Failed to pull task - will retry on stream restart");
1714 return true; } else {
1716 *pull_task_ack = false;
1717 }
1718 }
1719 }
1720
1721 let running_count = task_queue.running_parallelism_sum();
1722 let waiting_count = task_queue.waiting_parallelism_sum();
1723 let available_count = max_task_parallelism.saturating_sub(running_count);
1724 let current_state = IcebergCompactionLogState {
1725 running_parallelism: running_count,
1726 waiting_parallelism: waiting_count,
1727 available_parallelism: available_count,
1728 pull_task_ack: *pull_task_ack,
1729 pending_pull_task_count,
1730 };
1731
1732 if log_throttler.should_log(¤t_state) {
1734 tracing::info!(
1735 running_parallelism_count = %current_state.running_parallelism,
1736 waiting_parallelism_count = %current_state.waiting_parallelism,
1737 available_parallelism = %current_state.available_parallelism,
1738 pull_task_ack = %current_state.pull_task_ack,
1739 pending_pull_task_count = %current_state.pending_pull_task_count
1740 );
1741 log_throttler.update(current_state);
1742 }
1743
1744 false }
1746
1747#[cfg(test)]
1748mod tests {
1749 use super::*;
1750
1751 #[test]
1752 fn test_cancel_iceberg_task_removes_waiting_plans_and_tracker() {
1753 let task_id = IcebergCompactionTaskId::new(42);
1754 let mut task_queue = IcebergTaskQueue::new(10, 30, 1024);
1755 assert_eq!(
1756 task_queue.push(
1757 iceberg_compaction::IcebergTaskMeta {
1758 task_id,
1759 plan_index: 0,
1760 required_parallelism: 3,
1761 memory_reservation_bytes: 64,
1762 },
1763 None,
1764 ),
1765 PushResult::Added
1766 );
1767 assert_eq!(
1768 task_queue.push(
1769 iceberg_compaction::IcebergTaskMeta {
1770 task_id,
1771 plan_index: 1,
1772 required_parallelism: 4,
1773 memory_reservation_bytes: 64,
1774 },
1775 None,
1776 ),
1777 PushResult::Added
1778 );
1779 assert_eq!(task_queue.waiting_parallelism_sum(), 7);
1780
1781 let shutdown_map = Arc::new(Mutex::new(HashMap::new()));
1782 let mut task_trackers = HashMap::from([(task_id, IcebergTaskTracker::new(10, 2, false))]);
1783
1784 cancel_iceberg_task(task_id, &mut task_queue, &shutdown_map, &mut task_trackers);
1785
1786 assert_eq!(task_queue.waiting_parallelism_sum(), 0);
1787 assert!(!task_trackers.contains_key(&task_id));
1788 }
1789
1790 #[test]
1791 fn test_cancel_iceberg_task_shuts_down_running_plans_and_tracker() {
1792 let task_id = IcebergCompactionTaskId::new(43);
1793 let task_key = (task_id, 0);
1794 let mut task_queue = IcebergTaskQueue::new(10, 30, 1024);
1795 let shutdown_map = Arc::new(Mutex::new(HashMap::new()));
1796 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
1797 shutdown_map.lock().unwrap().insert(task_key, shutdown_tx);
1798 let mut task_trackers = HashMap::from([(task_id, IcebergTaskTracker::new(10, 1, false))]);
1799
1800 cancel_iceberg_task(task_id, &mut task_queue, &shutdown_map, &mut task_trackers);
1801
1802 assert!(shutdown_rx.try_recv().is_ok());
1803 assert!(shutdown_map.lock().unwrap().is_empty());
1804 assert!(!task_trackers.contains_key(&task_id));
1805 }
1806
1807 #[test]
1808 fn test_log_state_equality() {
1809 let state1 = CompactionLogState {
1811 running_parallelism: 10,
1812 pull_task_ack: true,
1813 pending_pull_task_count: 2,
1814 };
1815 let state2 = CompactionLogState {
1816 running_parallelism: 10,
1817 pull_task_ack: true,
1818 pending_pull_task_count: 2,
1819 };
1820 let state3 = CompactionLogState {
1821 running_parallelism: 11,
1822 pull_task_ack: true,
1823 pending_pull_task_count: 2,
1824 };
1825 assert_eq!(state1, state2);
1826 assert_ne!(state1, state3);
1827
1828 let ice_state1 = IcebergCompactionLogState {
1830 running_parallelism: 10,
1831 waiting_parallelism: 5,
1832 available_parallelism: 15,
1833 pull_task_ack: true,
1834 pending_pull_task_count: 2,
1835 };
1836 let ice_state2 = IcebergCompactionLogState {
1837 running_parallelism: 10,
1838 waiting_parallelism: 6,
1839 available_parallelism: 15,
1840 pull_task_ack: true,
1841 pending_pull_task_count: 2,
1842 };
1843 assert_ne!(ice_state1, ice_state2);
1844 }
1845
1846 #[test]
1847 fn test_log_throttler_state_change_detection() {
1848 let mut throttler = LogThrottler::<CompactionLogState>::new(Duration::from_secs(60));
1849 let state1 = CompactionLogState {
1850 running_parallelism: 10,
1851 pull_task_ack: true,
1852 pending_pull_task_count: 2,
1853 };
1854 let state2 = CompactionLogState {
1855 running_parallelism: 11,
1856 pull_task_ack: true,
1857 pending_pull_task_count: 2,
1858 };
1859
1860 assert!(throttler.should_log(&state1));
1862 throttler.update(state1.clone());
1863
1864 assert!(!throttler.should_log(&state1));
1866
1867 assert!(throttler.should_log(&state2));
1869 throttler.update(state2.clone());
1870
1871 assert!(!throttler.should_log(&state2));
1873 }
1874
1875 #[test]
1876 fn test_log_throttler_heartbeat() {
1877 let mut throttler = LogThrottler::<CompactionLogState>::new(Duration::from_millis(10));
1878 let state = CompactionLogState {
1879 running_parallelism: 10,
1880 pull_task_ack: true,
1881 pending_pull_task_count: 2,
1882 };
1883
1884 assert!(throttler.should_log(&state));
1886 throttler.update(state.clone());
1887
1888 assert!(!throttler.should_log(&state));
1890
1891 std::thread::sleep(Duration::from_millis(15));
1893
1894 assert!(throttler.should_log(&state));
1896 }
1897}