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