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