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