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