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_iceberg_task_report, create_task_execution,
67    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) -> (JoinHandle<()>, Sender<()>) {
403    let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
404    let stream_retry_interval = Duration::from_secs(30);
405    let periodic_event_update_interval = Duration::from_millis(
406        compactor_context
407            .storage_opts
408            .iceberg_compaction_pull_interval_ms,
409    );
410    let worker_num = compactor_context.compaction_executor.worker_num();
411
412    let max_task_parallelism: u32 = (worker_num as f32
413        * compactor_context.storage_opts.compactor_max_task_multiplier)
414        .ceil() as u32;
415
416    const MAX_PULL_TASK_COUNT: u32 = 4;
417    let max_pull_task_count = std::cmp::min(max_task_parallelism, MAX_PULL_TASK_COUNT);
418
419    assert_ge!(
420        compactor_context.storage_opts.compactor_max_task_multiplier,
421        0.0
422    );
423
424    let join_handle = tokio::spawn(async move {
425        // Initialize task queue with event-driven scheduling
426        let pending_parallelism_budget = (max_task_parallelism as f32
427            * compactor_context
428                .storage_opts
429                .iceberg_compaction_pending_parallelism_budget_multiplier)
430            .ceil() as u32;
431        let mut task_queue =
432            IcebergTaskQueue::new(max_task_parallelism, pending_parallelism_budget);
433
434        // Shutdown tracking for running tasks (task_key -> shutdown_sender)
435        let shutdown_map = Arc::new(Mutex::new(HashMap::<TaskKey, Sender<()>>::new()));
436
437        // Channel for task completion notifications
438        let (task_completion_tx, mut task_completion_rx) =
439            tokio::sync::mpsc::unbounded_channel::<IcebergPlanCompletion>();
440        let mut task_trackers = HashMap::<IcebergCompactionTaskId, IcebergTaskTracker>::new();
441        // Buffers task reports that failed to send on the current stream.
442        // The queue is flushed in FIFO order after the stream reconnects.
443        let mut pending_task_reports = VecDeque::<IcebergTaskReport>::new();
444
445        let mut min_interval = tokio::time::interval(stream_retry_interval);
446        let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
447
448        // Track last logged state to avoid duplicate logs
449        let mut log_throttler =
450            LogThrottler::<IcebergCompactionLogState>::new(COMPACTION_HEARTBEAT_LOG_INTERVAL);
451
452        // This outer loop is to recreate stream.
453        'start_stream: loop {
454            // reset state
455            // pull_task_ack.store(true, Ordering::SeqCst);
456            let mut pull_task_ack = true;
457            tokio::select! {
458                // Wait for interval.
459                _ = min_interval.tick() => {},
460                // Shutdown compactor.
461                _ = &mut shutdown_rx => {
462                    tracing::info!("Compactor is shutting down");
463                    return;
464                }
465            }
466
467            let (request_sender, response_event_stream) = match hummock_meta_client
468                .subscribe_iceberg_compaction_event()
469                .await
470            {
471                Ok((request_sender, response_event_stream)) => {
472                    tracing::debug!("Succeeded subscribe_iceberg_compaction_event.");
473                    (request_sender, response_event_stream)
474                }
475
476                Err(e) => {
477                    tracing::warn!(
478                        error = %e.as_report(),
479                        "Subscribing to iceberg compaction tasks failed with error. Will retry.",
480                    );
481                    continue 'start_stream;
482                }
483            };
484
485            if matches!(
486                flush_pending_iceberg_task_reports(&request_sender, &mut pending_task_reports),
487                ReportSendResult::RestartStream
488            ) {
489                continue 'start_stream;
490            }
491
492            pin_mut!(response_event_stream);
493
494            let _executor = compactor_context.compaction_executor.clone();
495
496            // This inner loop is to consume stream or report task progress.
497            let mut event_loop_iteration_now = Instant::now();
498            'consume_stream: loop {
499                {
500                    // report
501                    compactor_context
502                        .compactor_metrics
503                        .compaction_event_loop_iteration_latency
504                        .observe(event_loop_iteration_now.elapsed().as_millis() as _);
505                    event_loop_iteration_now = Instant::now();
506                }
507
508                let request_sender = request_sender.clone();
509                let event: Option<Result<SubscribeIcebergCompactionEventResponse, _>> = tokio::select! {
510                    // Handle task completion notifications
511                    Some(plan_completion) = task_completion_rx.recv() => {
512                        let task_key = plan_completion.task_key;
513                        let error_message = plan_completion.error_message;
514                        let pk_index_result = plan_completion.pk_index_result;
515                        tracing::debug!(
516                            task_id = %task_key.0,
517                            plan_index = task_key.1,
518                            success = error_message.is_none(),
519                            "Plan completed, updating queue state"
520                        );
521                        task_queue.finish_running(task_key);
522
523                        let completed_task_id = task_key.0;
524                        let Entry::Occupied(mut tracker_entry) =
525                            task_trackers.entry(completed_task_id)
526                        else {
527                            continue 'consume_stream;
528                        };
529                        tracker_entry
530                            .get_mut()
531                            .record_completion(error_message, pk_index_result);
532                        if !tracker_entry.get().is_finished() {
533                            continue 'consume_stream;
534                        }
535
536                        let tracker = tracker_entry.remove();
537                        let sink_id = tracker.sink_id();
538                        let total_plans = tracker.total_plans();
539                        let successful_plans = tracker.successful_plans();
540                        let failed_plans = tracker.failed_plans();
541                        let report = tracker.into_report(completed_task_id);
542                        let task_succeeded = report.error_message.is_none();
543                        if matches!(
544                            send_or_buffer_iceberg_task_report(
545                                &request_sender,
546                                &mut pending_task_reports,
547                                report,
548                            ),
549                            ReportSendResult::RestartStream
550                        ) {
551                            continue 'start_stream;
552                        }
553                        if task_succeeded {
554                            tracing::info!(
555                                iceberg_component = "compaction_worker",
556                                iceberg_operation = "report_task",
557                                task_id = %completed_task_id,
558                                sink_id = sink_id,
559                                total_plans = total_plans,
560                                successful_plans = successful_plans,
561                                failed_plans = failed_plans,
562                                "iceberg_compaction_task_succeeded",
563                            );
564                        }
565                        continue 'consume_stream;
566                    }
567
568                    // Event-driven task scheduling - wait for tasks to become schedulable
569                    _ = task_queue.wait_schedulable() => {
570                        schedule_queued_tasks(
571                            &mut task_queue,
572                            &compactor_context,
573                            &shutdown_map,
574                            &task_completion_tx,
575                        );
576                        continue 'consume_stream;
577                    }
578
579                    _ = periodic_event_interval.tick() => {
580                        // Only handle meta task pulling in periodic tick
581                        let should_restart_stream = handle_meta_task_pulling(
582                            &mut pull_task_ack,
583                            &task_queue,
584                            max_task_parallelism,
585                            max_pull_task_count,
586                            &request_sender,
587                            &mut log_throttler,
588                        );
589
590                        if should_restart_stream {
591                            continue 'start_stream;
592                        }
593                        continue;
594                    }
595                    event = response_event_stream.next() => {
596                        event
597                    }
598
599                    _ = &mut shutdown_rx => {
600                        tracing::info!("Iceberg Compactor is shutting down");
601                        return
602                    }
603                };
604
605                match event {
606                    Some(Ok(SubscribeIcebergCompactionEventResponse {
607                        event,
608                        create_at: _create_at,
609                    })) => {
610                        let event = match event {
611                            Some(event) => event,
612                            None => continue 'consume_stream,
613                        };
614
615                        match event {
616                            risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::CompactTask(iceberg_compaction_task) => {
617                                let task_id = iceberg_compaction_task.task_id;
618                                let sink_id = iceberg_compaction_task.sink_id;
619                                // iceberg-rust internally spawns on native Tokio, which is not
620                                // available inside a madsim node. Tests can replace only the
621                                // rewrite result while retaining the real compactor stream and
622                                // all downstream resolver/commit/recovery behavior.
623                                #[cfg(madsim)]
624                                if iceberg_compaction_task.pk_index_coordinated
625                                    && fail::eval(
626                                        "iceberg_pk_index_simulated_compaction_report",
627                                        |_| (),
628                                    )
629                                    .is_some()
630                                    && let Some(result) = simulated_pk_index_compaction_result()
631                                {
632                                    let mut report = build_iceberg_task_report(task_id, sink_id, None);
633                                    report.pk_index_result = Some(result);
634                                    if matches!(
635                                        send_or_buffer_iceberg_task_report(
636                                            &request_sender,
637                                            &mut pending_task_reports,
638                                            report,
639                                        ),
640                                        ReportSendResult::RestartStream
641                                    ) {
642                                        continue 'start_stream;
643                                    }
644                                    continue 'consume_stream;
645                                }
646                                // Note: write_parquet_properties is now built from sink config (IcebergConfig) in create_task_execution
647                                let compactor_runner_config = match IcebergCompactorRunnerConfigBuilder::default()
648                                    .max_parallelism((worker_num as f32 * compactor_context.storage_opts.iceberg_compaction_task_parallelism_ratio) as u32)
649                                    .min_size_per_partition(compactor_context.storage_opts.iceberg_compaction_min_size_per_partition_mb as u64 * 1024 * 1024)
650                                    .max_file_count_per_partition(compactor_context.storage_opts.iceberg_compaction_max_file_count_per_partition)
651                                    .enable_validate_compaction(compactor_context.storage_opts.iceberg_compaction_enable_validate)
652                                    .max_record_batch_rows(compactor_context.storage_opts.iceberg_compaction_max_record_batch_rows)
653                                    .enable_heuristic_output_parallelism(compactor_context.storage_opts.iceberg_compaction_enable_heuristic_output_parallelism)
654                                    .max_concurrent_closes(compactor_context.storage_opts.iceberg_compaction_max_concurrent_closes)
655                                    .enable_prefetch(compactor_context.storage_opts.iceberg_compaction_enable_prefetch)
656                                    .target_binpack_group_size_mb(
657                                        compactor_context.storage_opts.iceberg_compaction_target_binpack_group_size_mb
658                                    )
659                                    .min_group_size_mb(
660                                        compactor_context.storage_opts.iceberg_compaction_min_group_size_mb
661                                    )
662                                    .min_group_file_count(
663                                        compactor_context.storage_opts.iceberg_compaction_min_group_file_count
664                                    )
665                                    .build() {
666                                    Ok(config) => config,
667                                    Err(e) => {
668                                        tracing::warn!(
669                                            iceberg_component = "compaction_worker",
670                                            iceberg_operation = "build_runner_config",
671                                            error = %e.as_report(),
672                                            task_id = %task_id,
673                                            sink_id = sink_id,
674                                            "iceberg_compaction_runner_config_failed",
675                                        );
676                                        let report = build_iceberg_task_report(
677                                            task_id,
678                                            sink_id,
679                                            Some(format!(
680                                                "Failed to build iceberg compactor runner config: {}",
681                                                e.as_report()
682                                            )),
683                                        );
684                                        if matches!(
685                                            send_or_buffer_iceberg_task_report(
686                                                &request_sender,
687                                                &mut pending_task_reports,
688                                                report,
689                                            ),
690                                            ReportSendResult::RestartStream
691                                        ) {
692                                            continue 'start_stream;
693                                        }
694                                        continue 'consume_stream;
695                                    }
696                                };
697
698                                // Create task execution context and plan runners from the task
699                                let task_execution = match create_task_execution(
700                                    iceberg_compaction_task,
701                                    compactor_runner_config,
702                                    compactor_context.compactor_metrics.clone(),
703                                ).await {
704                                    Ok(task_execution) => task_execution,
705                                    Err(e) => {
706                                        tracing::warn!(
707                                            iceberg_component = "compaction_worker",
708                                            iceberg_operation = "plan_task",
709                                            error = %e.as_report(),
710                                            task_id = %task_id,
711                                            sink_id = sink_id,
712                                            "iceberg_compaction_task_plan_failed",
713                                        );
714                                        let report = build_iceberg_task_report(
715                                            task_id,
716                                            sink_id,
717                                            Some(format!(
718                                                "Failed to create iceberg compaction task execution: {}",
719                                                e.as_report()
720                                            )),
721                                        );
722                                        if matches!(
723                                            send_or_buffer_iceberg_task_report(
724                                                &request_sender,
725                                                &mut pending_task_reports,
726                                                report,
727                                            ),
728                                            ReportSendResult::RestartStream
729                                        ) {
730                                            continue 'start_stream;
731                                        }
732                                        continue 'consume_stream;
733                                    }
734                                };
735
736                                let sink_id = task_execution.sink_id;
737                                let plan_runners = task_execution.plan_runners;
738
739                                if plan_runners.is_empty() {
740                                    tracing::info!(
741                                        iceberg_component = "compaction_worker",
742                                        iceberg_operation = "enqueue_plan",
743                                        task_id = %task_id,
744                                        sink_id = sink_id,
745                                        "iceberg_compaction_task_skipped_no_plans",
746                                    );
747                                    let report = build_iceberg_task_report(task_id, sink_id, None);
748                                    if matches!(
749                                        send_or_buffer_iceberg_task_report(
750                                            &request_sender,
751                                            &mut pending_task_reports,
752                                            report,
753                                        ),
754                                        ReportSendResult::RestartStream
755                                    ) {
756                                        continue 'start_stream;
757                                    }
758                                    continue 'consume_stream;
759                                }
760
761                                // Enqueue each plan runner independently
762                                let total_plans = plan_runners.len();
763                                let mut enqueued_count = 0;
764
765                                for runner in plan_runners {
766                                    let meta = runner.to_meta();
767                                    let plan_index = meta.plan_index;
768                                    let required_parallelism = runner.required_parallelism();
769                                    let runner_sink_id = runner.sink_id;
770                                    let runner_task_type = runner.task_type;
771                                    let runner_table = runner.table_ident.to_string();
772                                    let push_result = task_queue.push(meta, Some(runner));
773
774                                    match push_result {
775                                        PushResult::Added => {
776                                            enqueued_count += 1;
777                                            tracing::debug!(
778                                                iceberg_component = "compaction_worker",
779                                                iceberg_operation = "enqueue_plan",
780                                                task_id = %task_id,
781                                                sink_id = runner_sink_id,
782                                                plan_index = plan_index,
783                                                task_type = ?runner_task_type,
784                                                table = %runner_table,
785                                                required_parallelism = required_parallelism,
786                                                "iceberg_compaction_plan_enqueued",
787                                            );
788                                        },
789                                        PushResult::RejectedCapacity => {
790                                            tracing::warn!(
791                                                iceberg_component = "compaction_worker",
792                                                iceberg_operation = "enqueue_plan",
793                                                task_id = %task_id,
794                                                sink_id = runner_sink_id,
795                                                plan_index = plan_index,
796                                                task_type = ?runner_task_type,
797                                                table = %runner_table,
798                                                required_parallelism = required_parallelism,
799                                                pending_budget = pending_parallelism_budget,
800                                                enqueued_count = enqueued_count,
801                                                total_plans = total_plans,
802                                                "iceberg_compaction_plan_rejected_capacity",
803                                            );
804                                            // Stop enqueuing remaining plans
805                                            break;
806                                        },
807                                        PushResult::RejectedTooLarge => {
808                                            tracing::error!(
809                                                iceberg_component = "compaction_worker",
810                                                iceberg_operation = "enqueue_plan",
811                                                task_id = %task_id,
812                                                sink_id = runner_sink_id,
813                                                plan_index = plan_index,
814                                                task_type = ?runner_task_type,
815                                                table = %runner_table,
816                                                required_parallelism = required_parallelism,
817                                                max_parallelism = max_task_parallelism,
818                                                "iceberg_compaction_plan_rejected_too_large",
819                                            );
820                                        },
821                                        PushResult::RejectedInvalidParallelism => {
822                                            tracing::error!(
823                                                iceberg_component = "compaction_worker",
824                                                iceberg_operation = "enqueue_plan",
825                                                task_id = %task_id,
826                                                sink_id = runner_sink_id,
827                                                plan_index = plan_index,
828                                                task_type = ?runner_task_type,
829                                                table = %runner_table,
830                                                required_parallelism = required_parallelism,
831                                                "iceberg_compaction_plan_rejected_invalid_parallelism",
832                                            );
833                                        },
834                                        PushResult::RejectedDuplicate => {
835                                            tracing::error!(
836                                                iceberg_component = "compaction_worker",
837                                                iceberg_operation = "enqueue_plan",
838                                                task_id = %task_id,
839                                                sink_id = runner_sink_id,
840                                                plan_index = plan_index,
841                                                task_type = ?runner_task_type,
842                                                table = %runner_table,
843                                                "iceberg_compaction_plan_rejected_duplicate",
844                                            );
845                                        }
846                                    }
847                                }
848
849                                if enqueued_count == 0 {
850                                    let report = build_iceberg_task_report(
851                                        task_id,
852                                        sink_id,
853                                        Some("Failed to enqueue all iceberg compaction plans".to_owned()),
854                                    );
855                                    if matches!(
856                                        send_or_buffer_iceberg_task_report(
857                                            &request_sender,
858                                            &mut pending_task_reports,
859                                            report,
860                                        ),
861                                        ReportSendResult::RestartStream
862                                    ) {
863                                        continue 'start_stream;
864                                    }
865                                } else {
866                                    task_trackers.insert(
867                                        task_id,
868                                        IcebergTaskTracker::new(sink_id, enqueued_count),
869                                    );
870                                }
871
872                                tracing::info!(
873                                    iceberg_component = "compaction_worker",
874                                    iceberg_operation = "enqueue_plan",
875                                    task_id = %task_id,
876                                    sink_id = sink_id,
877                                    total_plans = total_plans,
878                                    enqueued_count = enqueued_count,
879                                    "iceberg_compaction_task_enqueue_finished",
880                                );
881                            },
882                            risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::PullTaskAck(_) => {
883                                // set flag
884                                pull_task_ack = true;
885                            },
886                            risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::Event::CancelCompactTask(cancel_compact_task) => {
887                                cancel_iceberg_task(
888                                    cancel_compact_task.task_id,
889                                    &mut task_queue,
890                                    &shutdown_map,
891                                    &mut task_trackers,
892                                );
893                            },
894                        }
895                    }
896                    Some(Err(e)) => {
897                        tracing::warn!("Failed to consume stream. {}", e.message());
898                        continue 'start_stream;
899                    }
900                    _ => {
901                        // The stream is exhausted
902                        continue 'start_stream;
903                    }
904                }
905            }
906        }
907    });
908
909    (join_handle, shutdown_tx)
910}
911
912/// The background compaction thread that receives compaction tasks from hummock compaction
913/// manager and runs compaction tasks.
914#[must_use]
915pub fn start_compactor(
916    compactor_context: CompactorContext,
917    hummock_meta_client: Arc<dyn HummockMetaClient>,
918    object_id_manager: Arc<ObjectIdManager>,
919    compaction_catalog_manager_ref: CompactionCatalogManagerRef,
920) -> (JoinHandle<()>, Sender<()>) {
921    type CompactionShutdownMap = Arc<Mutex<HashMap<u64, Sender<()>>>>;
922    let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
923    let stream_retry_interval = Duration::from_secs(30);
924    let task_progress = compactor_context.task_progress_manager.clone();
925    let periodic_event_update_interval = Duration::from_millis(1000);
926
927    let max_task_parallelism: u32 = (compactor_context.compaction_executor.worker_num() as f32
928        * compactor_context.storage_opts.compactor_max_task_multiplier)
929        .ceil() as u32;
930    let running_task_parallelism = Arc::new(AtomicU32::new(0));
931
932    const MAX_PULL_TASK_COUNT: u32 = 4;
933    let max_pull_task_count = std::cmp::min(max_task_parallelism, MAX_PULL_TASK_COUNT);
934
935    assert_ge!(
936        compactor_context.storage_opts.compactor_max_task_multiplier,
937        0.0
938    );
939
940    let join_handle = tokio::spawn(async move {
941        let shutdown_map = CompactionShutdownMap::default();
942        let mut min_interval = tokio::time::interval(stream_retry_interval);
943        let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
944
945        // Track last logged state to avoid duplicate logs
946        let mut log_throttler =
947            LogThrottler::<CompactionLogState>::new(COMPACTION_HEARTBEAT_LOG_INTERVAL);
948
949        // This outer loop is to recreate stream.
950        'start_stream: loop {
951            // reset state
952            // pull_task_ack.store(true, Ordering::SeqCst);
953            let mut pull_task_ack = true;
954            tokio::select! {
955                // Wait for interval.
956                _ = min_interval.tick() => {},
957                // Shutdown compactor.
958                _ = &mut shutdown_rx => {
959                    tracing::info!("Compactor is shutting down");
960                    return;
961                }
962            }
963
964            let (request_sender, response_event_stream) =
965                match hummock_meta_client.subscribe_compaction_event().await {
966                    Ok((request_sender, response_event_stream)) => {
967                        tracing::debug!("Succeeded subscribe_compaction_event.");
968                        (request_sender, response_event_stream)
969                    }
970
971                    Err(e) => {
972                        tracing::warn!(
973                            error = %e.as_report(),
974                            "Subscribing to compaction tasks failed with error. Will retry.",
975                        );
976                        continue 'start_stream;
977                    }
978                };
979
980            pin_mut!(response_event_stream);
981
982            let executor = compactor_context.compaction_executor.clone();
983            let object_id_manager = object_id_manager.clone();
984
985            // This inner loop is to consume stream or report task progress.
986            let mut event_loop_iteration_now = Instant::now();
987            'consume_stream: loop {
988                {
989                    // report
990                    compactor_context
991                        .compactor_metrics
992                        .compaction_event_loop_iteration_latency
993                        .observe(event_loop_iteration_now.elapsed().as_millis() as _);
994                    event_loop_iteration_now = Instant::now();
995                }
996
997                let running_task_parallelism = running_task_parallelism.clone();
998                let request_sender = request_sender.clone();
999                let event: Option<Result<SubscribeCompactionEventResponse, _>> = tokio::select! {
1000                    _ = periodic_event_interval.tick() => {
1001                        let progress_list = get_task_progress(task_progress.clone());
1002
1003                        if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1004                            event: Some(RequestEvent::HeartBeat(
1005                                HeartBeat {
1006                                    progress: progress_list
1007                                }
1008                            )),
1009                            create_at: SystemTime::now()
1010                                .duration_since(std::time::UNIX_EPOCH)
1011                                .expect("Clock may have gone backwards")
1012                                .as_millis() as u64,
1013                        }) {
1014                            tracing::warn!(error = %e.as_report(), "Failed to report task progress");
1015                            // re subscribe stream
1016                            continue 'start_stream;
1017                        }
1018
1019
1020                        let mut pending_pull_task_count = 0;
1021                        if pull_task_ack {
1022                            // TODO: Compute parallelism on meta side
1023                            pending_pull_task_count = (max_task_parallelism - running_task_parallelism.load(Ordering::SeqCst)).min(max_pull_task_count);
1024
1025                            if pending_pull_task_count > 0 {
1026                                if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1027                                    event: Some(RequestEvent::PullTask(
1028                                        PullTask {
1029                                            pull_task_count: pending_pull_task_count,
1030                                        }
1031                                    )),
1032                                    create_at: SystemTime::now()
1033                                        .duration_since(std::time::UNIX_EPOCH)
1034                                        .expect("Clock may have gone backwards")
1035                                        .as_millis() as u64,
1036                                }) {
1037                                    tracing::warn!(error = %e.as_report(), "Failed to pull task");
1038
1039                                    // re subscribe stream
1040                                    continue 'start_stream;
1041                                } else {
1042                                    pull_task_ack = false;
1043                                }
1044                            }
1045                        }
1046
1047                        let running_count = running_task_parallelism.load(Ordering::SeqCst);
1048                        let current_state = CompactionLogState {
1049                            running_parallelism: running_count,
1050                            pull_task_ack,
1051                            pending_pull_task_count,
1052                        };
1053
1054                        // Log only when state changes or periodically as heartbeat
1055                        if log_throttler.should_log(&current_state) {
1056                            tracing::info!(
1057                                running_parallelism_count = %current_state.running_parallelism,
1058                                pull_task_ack = %current_state.pull_task_ack,
1059                                pending_pull_task_count = %current_state.pending_pull_task_count
1060                            );
1061                            log_throttler.update(current_state);
1062                        }
1063
1064                        continue;
1065                    }
1066                    event = response_event_stream.next() => {
1067                        event
1068                    }
1069
1070                    _ = &mut shutdown_rx => {
1071                        tracing::info!("Compactor is shutting down");
1072                        return
1073                    }
1074                };
1075
1076                fn send_report_task_event(
1077                    compact_task: &CompactTask,
1078                    table_stats: TableStatsMap,
1079                    object_timestamps: HashMap<HummockSstableObjectId, u64>,
1080                    request_sender: &mpsc::UnboundedSender<SubscribeCompactionEventRequest>,
1081                ) {
1082                    if let Err(e) = request_sender.send(SubscribeCompactionEventRequest {
1083                        event: Some(RequestEvent::ReportTask(ReportTask {
1084                            task_id: compact_task.task_id,
1085                            task_status: compact_task.task_status.into(),
1086                            sorted_output_ssts: compact_task
1087                                .sorted_output_ssts
1088                                .iter()
1089                                .map(|sst| sst.into())
1090                                .collect(),
1091                            table_stats_change: to_prost_table_stats_map(table_stats),
1092                            object_timestamps,
1093                        })),
1094                        create_at: SystemTime::now()
1095                            .duration_since(std::time::UNIX_EPOCH)
1096                            .expect("Clock may have gone backwards")
1097                            .as_millis() as u64,
1098                    }) {
1099                        let task_id = compact_task.task_id;
1100                        tracing::warn!(error = %e.as_report(), "Failed to report task {task_id:?}");
1101                    }
1102                }
1103
1104                match event {
1105                    Some(Ok(SubscribeCompactionEventResponse { event, create_at })) => {
1106                        let event = match event {
1107                            Some(event) => event,
1108                            None => continue 'consume_stream,
1109                        };
1110                        let shutdown = shutdown_map.clone();
1111                        let context = compactor_context.clone();
1112                        let consumed_latency_ms = SystemTime::now()
1113                            .duration_since(std::time::UNIX_EPOCH)
1114                            .expect("Clock may have gone backwards")
1115                            .as_millis() as u64
1116                            - create_at;
1117                        context
1118                            .compactor_metrics
1119                            .compaction_event_consumed_latency
1120                            .observe(consumed_latency_ms as _);
1121
1122                        let object_id_manager = object_id_manager.clone();
1123                        let compaction_catalog_manager_ref = compaction_catalog_manager_ref.clone();
1124
1125                        match event {
1126                            ResponseEvent::CompactTask(compact_task) => {
1127                                let compact_task = CompactTask::from(compact_task);
1128                                let parallelism =
1129                                    calculate_task_parallelism(&compact_task, &context);
1130
1131                                assert_ne!(parallelism, 0, "splits cannot be empty");
1132
1133                                if (max_task_parallelism
1134                                    - running_task_parallelism.load(Ordering::SeqCst))
1135                                    < parallelism as u32
1136                                {
1137                                    tracing::warn!(
1138                                        "Not enough core parallelism to serve the task {} task_parallelism {} running_task_parallelism {} max_task_parallelism {}",
1139                                        compact_task.task_id,
1140                                        parallelism,
1141                                        max_task_parallelism,
1142                                        running_task_parallelism.load(Ordering::Relaxed),
1143                                    );
1144                                    let (compact_task, table_stats, object_timestamps) =
1145                                        compact_done(
1146                                            compact_task,
1147                                            context.clone(),
1148                                            vec![],
1149                                            TaskStatus::NoAvailCpuResourceCanceled,
1150                                        );
1151
1152                                    send_report_task_event(
1153                                        &compact_task,
1154                                        table_stats,
1155                                        object_timestamps,
1156                                        &request_sender,
1157                                    );
1158
1159                                    continue 'consume_stream;
1160                                }
1161
1162                                running_task_parallelism
1163                                    .fetch_add(parallelism as u32, Ordering::SeqCst);
1164                                executor.spawn(async move {
1165                                    let (tx, rx) = tokio::sync::oneshot::channel();
1166                                    let task_id = compact_task.task_id;
1167                                    shutdown.lock().unwrap().insert(task_id, tx);
1168
1169                                    let (
1170                                        (compact_task, table_stats, object_timestamps),
1171                                        _memory_tracker,
1172                                    ) = compactor_runner::compact(
1173                                        context.clone(),
1174                                        compact_task,
1175                                        rx,
1176                                        object_id_manager.clone(),
1177                                        compaction_catalog_manager_ref.clone(),
1178                                    )
1179                                    .await;
1180
1181                                    shutdown.lock().unwrap().remove(&task_id);
1182                                    running_task_parallelism
1183                                        .fetch_sub(parallelism as u32, Ordering::SeqCst);
1184
1185                                    send_report_task_event(
1186                                        &compact_task,
1187                                        table_stats,
1188                                        object_timestamps,
1189                                        &request_sender,
1190                                    );
1191
1192                                    let enable_check_compaction_result =
1193                                        context.storage_opts.check_compaction_result;
1194                                    let need_check_task =
1195                                        !compact_task.sorted_output_ssts.is_empty()
1196                                            && compact_task.task_status == TaskStatus::Success;
1197
1198                                    if enable_check_compaction_result && need_check_task {
1199                                        let read_table_ids = compact_task
1200                                            .get_table_ids_from_input_ssts()
1201                                            .collect::<Vec<_>>();
1202                                        match compaction_catalog_manager_ref.acquire(read_table_ids).await {
1203                                            Ok(compaction_catalog_agent_ref) =>  {
1204                                                match check_compaction_result(&compact_task, context.clone(), compaction_catalog_agent_ref).await
1205                                                {
1206                                                    Err(e) => {
1207                                                        tracing::warn!(error = %e.as_report(), "Failed to check compaction task {}",compact_task.task_id);
1208                                                    }
1209                                                    Ok(true) => (),
1210                                                    Ok(false) => {
1211                                                        panic!("Failed to pass consistency check for result of compaction task:\n{:?}", compact_task_to_string(&compact_task));
1212                                                    }
1213                                                }
1214                                            },
1215                                            Err(e) => {
1216                                                tracing::warn!(error = %e.as_report(), "failed to acquire compaction catalog agent");
1217                                            }
1218                                        }
1219                                    }
1220                                });
1221                            }
1222                            #[expect(deprecated)]
1223                            ResponseEvent::VacuumTask(_) => {
1224                                unreachable!("unexpected vacuum task");
1225                            }
1226                            #[expect(deprecated)]
1227                            ResponseEvent::FullScanTask(_) => {
1228                                unreachable!("unexpected scan task");
1229                            }
1230                            #[expect(deprecated)]
1231                            ResponseEvent::ValidationTask(validation_task) => {
1232                                let validation_task = ValidationTask::from(validation_task);
1233                                executor.spawn(async move {
1234                                    validate_ssts(validation_task, context.sstable_store.clone())
1235                                        .await;
1236                                });
1237                            }
1238                            ResponseEvent::CancelCompactTask(cancel_compact_task) => match shutdown
1239                                .lock()
1240                                .unwrap()
1241                                .remove(&cancel_compact_task.task_id)
1242                            {
1243                                Some(tx) => {
1244                                    if tx.send(()).is_err() {
1245                                        tracing::warn!(
1246                                            "Cancellation of compaction task failed. task_id: {}",
1247                                            cancel_compact_task.task_id
1248                                        );
1249                                    }
1250                                }
1251                                _ => {
1252                                    tracing::warn!(
1253                                        "Attempting to cancel non-existent compaction task. task_id: {}",
1254                                        cancel_compact_task.task_id
1255                                    );
1256                                }
1257                            },
1258
1259                            ResponseEvent::PullTaskAck(_pull_task_ack) => {
1260                                // set flag
1261                                pull_task_ack = true;
1262                            }
1263                        }
1264                    }
1265                    Some(Err(e)) => {
1266                        tracing::warn!("Failed to consume stream. {}", e.message());
1267                        continue 'start_stream;
1268                    }
1269                    _ => {
1270                        // The stream is exhausted
1271                        continue 'start_stream;
1272                    }
1273                }
1274            }
1275        }
1276    });
1277
1278    (join_handle, shutdown_tx)
1279}
1280
1281/// The background compaction thread that receives compaction tasks from hummock compaction
1282/// manager and runs compaction tasks.
1283#[must_use]
1284pub fn start_shared_compactor(
1285    grpc_proxy_client: GrpcCompactorProxyClient,
1286    mut receiver: mpsc::UnboundedReceiver<Request<DispatchCompactionTaskRequest>>,
1287    context: CompactorContext,
1288) -> (JoinHandle<()>, Sender<()>) {
1289    type CompactionShutdownMap = Arc<Mutex<HashMap<u64, Sender<()>>>>;
1290    let task_progress = context.task_progress_manager.clone();
1291    let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
1292    let periodic_event_update_interval = Duration::from_millis(1000);
1293
1294    let join_handle = tokio::spawn(async move {
1295        let shutdown_map = CompactionShutdownMap::default();
1296
1297        let mut periodic_event_interval = tokio::time::interval(periodic_event_update_interval);
1298        let executor = context.compaction_executor.clone();
1299        let report_heartbeat_client = grpc_proxy_client.clone();
1300        'consume_stream: loop {
1301            let request: Option<Request<DispatchCompactionTaskRequest>> = tokio::select! {
1302                _ = periodic_event_interval.tick() => {
1303                    let progress_list = get_task_progress(task_progress.clone());
1304                    let report_compaction_task_request = ReportCompactionTaskRequest{
1305                        event: Some(ReportCompactionTaskEvent::HeartBeat(
1306                            SharedHeartBeat {
1307                                progress: progress_list
1308                            }
1309                        )),
1310                     };
1311                    if let Err(e) = report_heartbeat_client.report_compaction_task(report_compaction_task_request).await{
1312                        tracing::warn!(error = %e.as_report(), "Failed to report heartbeat");
1313                    }
1314                    continue
1315                }
1316
1317
1318                _ = &mut shutdown_rx => {
1319                    tracing::info!("Compactor is shutting down");
1320                    return
1321                }
1322
1323                request = receiver.recv() => {
1324                    request
1325                }
1326
1327            };
1328            match request {
1329                Some(request) => {
1330                    let context = context.clone();
1331                    let shutdown = shutdown_map.clone();
1332
1333                    let cloned_grpc_proxy_client = grpc_proxy_client.clone();
1334                    executor.spawn(async move {
1335                        let DispatchCompactionTaskRequest {
1336                            tables,
1337                            output_object_ids,
1338                            task: dispatch_task,
1339                        } = request.into_inner();
1340                        let table_id_to_catalog = tables.into_iter().fold(HashMap::new(), |mut acc, table| {
1341                            acc.insert(table.id, table);
1342                            acc
1343                        });
1344
1345                        let mut output_object_ids_deque: VecDeque<_> = VecDeque::new();
1346                        output_object_ids_deque.extend(output_object_ids.into_iter().map(Into::<HummockSstableObjectId>::into));
1347                        let shared_compactor_object_id_manager =
1348                            SharedComapctorObjectIdManager::new(output_object_ids_deque, cloned_grpc_proxy_client.clone(), context.storage_opts.sstable_id_remote_fetch_number);
1349                            match dispatch_task.unwrap() {
1350                                dispatch_compaction_task_request::Task::CompactTask(compact_task) => {
1351                                    let compact_task = CompactTask::from(&compact_task);
1352                                    let (tx, rx) = tokio::sync::oneshot::channel();
1353                                    let task_id = compact_task.task_id;
1354                                    shutdown.lock().unwrap().insert(task_id, tx);
1355
1356                                    let compaction_catalog_manager_ref =
1357                                        Arc::new(CompactionCatalogManager::new_preloaded(
1358                                            table_id_to_catalog,
1359                                        ));
1360                                    let ((compact_task, table_stats, object_timestamps), _memory_tracker)= compactor_runner::compact(
1361                                        context.clone(),
1362                                        compact_task,
1363                                        rx,
1364                                        shared_compactor_object_id_manager,
1365                                        compaction_catalog_manager_ref.clone(),
1366                                    )
1367                                    .await;
1368                                    shutdown.lock().unwrap().remove(&task_id);
1369                                    let report_compaction_task_request = ReportCompactionTaskRequest {
1370                                        event: Some(ReportCompactionTaskEvent::ReportTask(ReportSharedTask {
1371                                            compact_task: Some(PbCompactTask::from(&compact_task)),
1372                                            table_stats_change: to_prost_table_stats_map(table_stats),
1373                                            object_timestamps,
1374                                    })),
1375                                    };
1376
1377                                    match cloned_grpc_proxy_client
1378                                        .report_compaction_task(report_compaction_task_request)
1379                                        .await
1380                                    {
1381                                        Ok(_) => {
1382                                            // TODO: remove this method after we have running risingwave cluster with fast compact algorithm stably for a long time.
1383                                            let enable_check_compaction_result =
1384                                                context.storage_opts.check_compaction_result;
1385                                            let need_check_task =
1386                                                !compact_task.sorted_output_ssts.is_empty()
1387                                                    && compact_task.task_status
1388                                                        == TaskStatus::Success;
1389                                            if enable_check_compaction_result && need_check_task {
1390                                                let read_table_ids = compact_task
1391                                                    .get_table_ids_from_input_ssts()
1392                                                    .collect::<Vec<_>>();
1393                                                match compaction_catalog_manager_ref.acquire(read_table_ids).await {
1394                                                    Ok(compaction_catalog_agent_ref) => {
1395                                                        match check_compaction_result(&compact_task, context.clone(), compaction_catalog_agent_ref).await
1396                                                        {
1397                                                            Err(e) => {
1398                                                                tracing::warn!(error = %e.as_report(), "Failed to check compaction task {}", task_id);
1399                                                            }
1400                                                            Ok(true) => (),
1401                                                            Ok(false) => {
1402                                                                panic!("Failed to pass consistency check for result of compaction task:\n{:?}", compact_task_to_string(&compact_task));
1403                                                            }
1404                                                        }
1405                                                    }
1406                                                    Err(e) => {
1407                                                        tracing::warn!(error = %e.as_report(), "failed to acquire compaction catalog agent");
1408                                                    }
1409                                                }
1410                                            }
1411                                        }
1412                                        Err(e) => tracing::warn!(error = %e.as_report(), "Failed to report task {task_id:?}"),
1413                                    }
1414
1415                                }
1416                                dispatch_compaction_task_request::Task::VacuumTask(_) => {
1417                                    unreachable!("unexpected vacuum task");
1418                                }
1419                                dispatch_compaction_task_request::Task::FullScanTask(_) => {
1420                                    unreachable!("unexpected scan task");
1421                                }
1422                                dispatch_compaction_task_request::Task::ValidationTask(validation_task) => {
1423                                    let validation_task = ValidationTask::from(validation_task);
1424                                    validate_ssts(validation_task, context.sstable_store.clone()).await;
1425                                }
1426                                dispatch_compaction_task_request::Task::CancelCompactTask(cancel_compact_task) => {
1427                                    match shutdown
1428                                        .lock()
1429                                        .unwrap()
1430                                        .remove(&cancel_compact_task.task_id)
1431                                    { Some(tx) => {
1432                                        if tx.send(()).is_err() {
1433                                            tracing::warn!(
1434                                                "Cancellation of compaction task failed. task_id: {}",
1435                                                cancel_compact_task.task_id
1436                                            );
1437                                        }
1438                                    } _ => {
1439                                        tracing::warn!(
1440                                            "Attempting to cancel non-existent compaction task. task_id: {}",
1441                                            cancel_compact_task.task_id
1442                                        );
1443                                    }}
1444                                }
1445                            }
1446                    });
1447                }
1448                None => continue 'consume_stream,
1449            }
1450        }
1451    });
1452    (join_handle, shutdown_tx)
1453}
1454
1455fn get_task_progress(
1456    task_progress: Arc<
1457        parking_lot::lock_api::Mutex<parking_lot::RawMutex, HashMap<u64, Arc<TaskProgress>>>,
1458    >,
1459) -> Vec<CompactTaskProgress> {
1460    let mut progress_list = Vec::new();
1461    for (&task_id, progress) in &*task_progress.lock() {
1462        progress_list.push(progress.snapshot(task_id));
1463    }
1464    progress_list
1465}
1466
1467/// Schedule queued tasks if we have capacity
1468fn schedule_queued_tasks(
1469    task_queue: &mut IcebergTaskQueue,
1470    compactor_context: &CompactorContext,
1471    shutdown_map: &Arc<Mutex<HashMap<TaskKey, Sender<()>>>>,
1472    task_completion_tx: &tokio::sync::mpsc::UnboundedSender<IcebergPlanCompletion>,
1473) {
1474    while let Some(popped_task) = task_queue.pop() {
1475        let task_id = popped_task.meta.task_id;
1476        let plan_index = popped_task.meta.plan_index;
1477        let task_key = (task_id, plan_index);
1478
1479        // Get unique_ident before moving runner
1480        let unique_ident = popped_task.runner.as_ref().map(|r| r.unique_ident());
1481
1482        let Some(runner) = popped_task.runner else {
1483            tracing::error!(
1484                iceberg_component = "compaction_worker",
1485                iceberg_operation = "schedule_plan",
1486                task_id = %task_id,
1487                plan_index = plan_index,
1488                "iceberg_compaction_plan_missing_runner",
1489            );
1490            task_queue.finish_running(task_key);
1491            continue;
1492        };
1493        let runner_sink_id = runner.sink_id;
1494        let runner_task_type = runner.task_type;
1495        let runner_table = runner.table_ident.to_string();
1496
1497        let executor = compactor_context.compaction_executor.clone();
1498        let shutdown_map_clone = shutdown_map.clone();
1499        let completion_tx_clone = task_completion_tx.clone();
1500        let (tx, rx) = tokio::sync::oneshot::channel();
1501
1502        {
1503            let mut shutdown_guard = shutdown_map.lock().unwrap();
1504            shutdown_guard.insert(task_key, tx);
1505        }
1506
1507        tracing::info!(
1508            iceberg_component = "compaction_worker",
1509            iceberg_operation = "schedule_plan",
1510            task_id = %task_id,
1511            sink_id = runner_sink_id,
1512            plan_index = plan_index,
1513            task_type = ?runner_task_type,
1514            table = %runner_table,
1515            unique_ident = ?unique_ident,
1516            required_parallelism = popped_task.meta.required_parallelism,
1517            "iceberg_compaction_plan_started_from_queue",
1518        );
1519
1520        executor.spawn(async move {
1521            let _cleanup_guard = scopeguard::guard(shutdown_map_clone, move |shutdown_map| {
1522                let mut shutdown_guard = shutdown_map.lock().unwrap();
1523                shutdown_guard.remove(&task_key);
1524            });
1525
1526            let result = Box::pin(runner.compact(rx)).await;
1527
1528            let completion = match result {
1529                Ok((_stats, pk_index_result)) => IcebergPlanCompletion {
1530                    task_key,
1531                    error_message: None,
1532                    pk_index_result,
1533                },
1534                Err(e) => {
1535                    if is_cancelled_iceberg_compaction_error(&e) {
1536                        tracing::info!(
1537                            iceberg_component = "compaction_worker",
1538                            iceberg_operation = "execute_plan",
1539                            task_id = %task_key.0,
1540                            sink_id = runner_sink_id,
1541                            plan_index = task_key.1,
1542                            task_type = ?runner_task_type,
1543                            table = %runner_table,
1544                            "iceberg_compaction_plan_cancelled",
1545                        );
1546                    } else {
1547                        tracing::warn!(
1548                            iceberg_component = "compaction_worker",
1549                            iceberg_operation = "execute_plan",
1550                            error = %e.as_report(),
1551                            task_id = %task_key.0,
1552                            sink_id = runner_sink_id,
1553                            plan_index = task_key.1,
1554                            task_type = ?runner_task_type,
1555                            table = %runner_table,
1556                            "iceberg_compaction_plan_failed",
1557                        );
1558                    }
1559                    IcebergPlanCompletion {
1560                        task_key,
1561                        error_message: Some(e.to_report_string()),
1562                        pk_index_result: None,
1563                    }
1564                }
1565            };
1566
1567            if completion_tx_clone.send(completion).is_err() {
1568                tracing::warn!(
1569                    iceberg_component = "compaction_worker",
1570                    iceberg_operation = "notify_plan_completion",
1571                    task_id = %task_key.0,
1572                    sink_id = runner_sink_id,
1573                    plan_index = task_key.1,
1574                    task_type = ?runner_task_type,
1575                    table = %runner_table,
1576                    "iceberg_compaction_plan_completion_send_failed",
1577                );
1578            }
1579        });
1580    }
1581}
1582
1583fn is_cancelled_iceberg_compaction_error(error: &crate::hummock::HummockError) -> bool {
1584    matches!(
1585        error.inner(),
1586        HummockErrorInner::CompactionExecutor(message) if message == "Plan cancelled"
1587    )
1588}
1589
1590fn cancel_iceberg_task(
1591    task_id: IcebergCompactionTaskId,
1592    task_queue: &mut IcebergTaskQueue,
1593    shutdown_map: &Arc<Mutex<HashMap<TaskKey, Sender<()>>>>,
1594    task_trackers: &mut HashMap<IcebergCompactionTaskId, IcebergTaskTracker>,
1595) {
1596    // Meta assigns one task id to an Iceberg compact task, but the compactor
1597    // splits it into multiple plan runners tracked by `(task_id, plan_index)`.
1598    // A cancel event only carries `task_id`, so it must cancel all waiting and
1599    // running plan runners that belong to that task.
1600    let cancelled_waiting = task_queue.cancel_waiting_task(task_id);
1601    let removed_tracker = task_trackers.remove(&task_id).is_some();
1602
1603    let cancelled_running = {
1604        let mut shutdown_guard = shutdown_map.lock().unwrap();
1605        let task_keys: Vec<_> = shutdown_guard
1606            .keys()
1607            .filter(|(running_task_id, _)| *running_task_id == task_id)
1608            .copied()
1609            .collect();
1610
1611        for task_key in &task_keys {
1612            if let Some(tx) = shutdown_guard.remove(task_key)
1613                && tx.send(()).is_err()
1614            {
1615                tracing::debug!(
1616                    task_id = %task_key.0,
1617                    plan_index = task_key.1,
1618                    "Iceberg compaction plan shutdown receiver already closed during cancellation"
1619                );
1620            }
1621        }
1622
1623        task_keys.len()
1624    };
1625
1626    if cancelled_waiting == 0 && cancelled_running == 0 && !removed_tracker {
1627        tracing::warn!(
1628            task_id = %task_id,
1629            "Attempting to cancel non-existent iceberg compaction task"
1630        );
1631    } else {
1632        tracing::info!(
1633            task_id = %task_id,
1634            cancelled_waiting = cancelled_waiting,
1635            cancelled_running = cancelled_running,
1636            removed_tracker = removed_tracker,
1637            "Cancelled iceberg compaction task"
1638        );
1639    }
1640}
1641
1642/// Handle pulling new tasks from meta service
1643/// Returns true if the stream should be restarted
1644fn handle_meta_task_pulling(
1645    pull_task_ack: &mut bool,
1646    task_queue: &IcebergTaskQueue,
1647    max_task_parallelism: u32,
1648    max_pull_task_count: u32,
1649    request_sender: &mpsc::UnboundedSender<SubscribeIcebergCompactionEventRequest>,
1650    log_throttler: &mut LogThrottler<IcebergCompactionLogState>,
1651) -> bool {
1652    let mut pending_pull_task_count = 0;
1653    if *pull_task_ack {
1654        // Use queue's running parallelism for pull decision
1655        let current_running_parallelism = task_queue.running_parallelism_sum();
1656        pending_pull_task_count =
1657            (max_task_parallelism - current_running_parallelism).min(max_pull_task_count);
1658
1659        if pending_pull_task_count > 0 {
1660            if let Err(e) = request_sender.send(SubscribeIcebergCompactionEventRequest {
1661                event: Some(subscribe_iceberg_compaction_event_request::Event::PullTask(
1662                    subscribe_iceberg_compaction_event_request::PullTask {
1663                        pull_task_count: pending_pull_task_count,
1664                    },
1665                )),
1666                create_at: SystemTime::now()
1667                    .duration_since(std::time::UNIX_EPOCH)
1668                    .expect("Clock may have gone backwards")
1669                    .as_millis() as u64,
1670            }) {
1671                tracing::warn!(error = %e.as_report(), "Failed to pull task - will retry on stream restart");
1672                return true; // Signal to restart stream
1673            } else {
1674                *pull_task_ack = false;
1675            }
1676        }
1677    }
1678
1679    let running_count = task_queue.running_parallelism_sum();
1680    let waiting_count = task_queue.waiting_parallelism_sum();
1681    let available_count = max_task_parallelism.saturating_sub(running_count);
1682    let current_state = IcebergCompactionLogState {
1683        running_parallelism: running_count,
1684        waiting_parallelism: waiting_count,
1685        available_parallelism: available_count,
1686        pull_task_ack: *pull_task_ack,
1687        pending_pull_task_count,
1688    };
1689
1690    // Log only when state changes or periodically as heartbeat
1691    if log_throttler.should_log(&current_state) {
1692        tracing::info!(
1693            running_parallelism_count = %current_state.running_parallelism,
1694            waiting_parallelism_count = %current_state.waiting_parallelism,
1695            available_parallelism = %current_state.available_parallelism,
1696            pull_task_ack = %current_state.pull_task_ack,
1697            pending_pull_task_count = %current_state.pending_pull_task_count
1698        );
1699        log_throttler.update(current_state);
1700    }
1701
1702    false // No need to restart stream
1703}
1704
1705#[cfg(test)]
1706mod tests {
1707    use super::*;
1708
1709    #[test]
1710    fn test_cancel_iceberg_task_removes_waiting_plans_and_tracker() {
1711        let task_id = IcebergCompactionTaskId::new(42);
1712        let mut task_queue = IcebergTaskQueue::new(10, 30);
1713        assert_eq!(
1714            task_queue.push(
1715                iceberg_compaction::IcebergTaskMeta {
1716                    task_id,
1717                    plan_index: 0,
1718                    required_parallelism: 3,
1719                },
1720                None,
1721            ),
1722            PushResult::Added
1723        );
1724        assert_eq!(
1725            task_queue.push(
1726                iceberg_compaction::IcebergTaskMeta {
1727                    task_id,
1728                    plan_index: 1,
1729                    required_parallelism: 4,
1730                },
1731                None,
1732            ),
1733            PushResult::Added
1734        );
1735        assert_eq!(task_queue.waiting_parallelism_sum(), 7);
1736
1737        let shutdown_map = Arc::new(Mutex::new(HashMap::new()));
1738        let mut task_trackers = HashMap::from([(task_id, IcebergTaskTracker::new(10, 2))]);
1739
1740        cancel_iceberg_task(task_id, &mut task_queue, &shutdown_map, &mut task_trackers);
1741
1742        assert_eq!(task_queue.waiting_parallelism_sum(), 0);
1743        assert!(!task_trackers.contains_key(&task_id));
1744    }
1745
1746    #[test]
1747    fn test_cancel_iceberg_task_shuts_down_running_plans_and_tracker() {
1748        let task_id = IcebergCompactionTaskId::new(43);
1749        let task_key = (task_id, 0);
1750        let mut task_queue = IcebergTaskQueue::new(10, 30);
1751        let shutdown_map = Arc::new(Mutex::new(HashMap::new()));
1752        let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
1753        shutdown_map.lock().unwrap().insert(task_key, shutdown_tx);
1754        let mut task_trackers = HashMap::from([(task_id, IcebergTaskTracker::new(10, 1))]);
1755
1756        cancel_iceberg_task(task_id, &mut task_queue, &shutdown_map, &mut task_trackers);
1757
1758        assert!(shutdown_rx.try_recv().is_ok());
1759        assert!(shutdown_map.lock().unwrap().is_empty());
1760        assert!(!task_trackers.contains_key(&task_id));
1761    }
1762
1763    #[test]
1764    fn test_log_state_equality() {
1765        // Test CompactionLogState
1766        let state1 = CompactionLogState {
1767            running_parallelism: 10,
1768            pull_task_ack: true,
1769            pending_pull_task_count: 2,
1770        };
1771        let state2 = CompactionLogState {
1772            running_parallelism: 10,
1773            pull_task_ack: true,
1774            pending_pull_task_count: 2,
1775        };
1776        let state3 = CompactionLogState {
1777            running_parallelism: 11,
1778            pull_task_ack: true,
1779            pending_pull_task_count: 2,
1780        };
1781        assert_eq!(state1, state2);
1782        assert_ne!(state1, state3);
1783
1784        // Test IcebergCompactionLogState
1785        let ice_state1 = IcebergCompactionLogState {
1786            running_parallelism: 10,
1787            waiting_parallelism: 5,
1788            available_parallelism: 15,
1789            pull_task_ack: true,
1790            pending_pull_task_count: 2,
1791        };
1792        let ice_state2 = IcebergCompactionLogState {
1793            running_parallelism: 10,
1794            waiting_parallelism: 6,
1795            available_parallelism: 15,
1796            pull_task_ack: true,
1797            pending_pull_task_count: 2,
1798        };
1799        assert_ne!(ice_state1, ice_state2);
1800    }
1801
1802    #[test]
1803    fn test_log_throttler_state_change_detection() {
1804        let mut throttler = LogThrottler::<CompactionLogState>::new(Duration::from_secs(60));
1805        let state1 = CompactionLogState {
1806            running_parallelism: 10,
1807            pull_task_ack: true,
1808            pending_pull_task_count: 2,
1809        };
1810        let state2 = CompactionLogState {
1811            running_parallelism: 11,
1812            pull_task_ack: true,
1813            pending_pull_task_count: 2,
1814        };
1815
1816        // First call should always log
1817        assert!(throttler.should_log(&state1));
1818        throttler.update(state1.clone());
1819
1820        // Same state should not log
1821        assert!(!throttler.should_log(&state1));
1822
1823        // Changed state should log
1824        assert!(throttler.should_log(&state2));
1825        throttler.update(state2.clone());
1826
1827        // Same state again should not log
1828        assert!(!throttler.should_log(&state2));
1829    }
1830
1831    #[test]
1832    fn test_log_throttler_heartbeat() {
1833        let mut throttler = LogThrottler::<CompactionLogState>::new(Duration::from_millis(10));
1834        let state = CompactionLogState {
1835            running_parallelism: 10,
1836            pull_task_ack: true,
1837            pending_pull_task_count: 2,
1838        };
1839
1840        // First call should log
1841        assert!(throttler.should_log(&state));
1842        throttler.update(state.clone());
1843
1844        // Same state immediately should not log
1845        assert!(!throttler.should_log(&state));
1846
1847        // Wait for heartbeat interval to pass
1848        std::thread::sleep(Duration::from_millis(15));
1849
1850        // Same state after interval should log (heartbeat)
1851        assert!(throttler.should_log(&state));
1852    }
1853}