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