Skip to main content

risingwave_storage/hummock/compactor/
mod.rs

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