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