1pub use self::iceberg_compactor_runner::create_task_execution;
16#[cfg(madsim)]
17pub use self::report::set_simulated_pk_index_compaction_result;
18#[cfg(madsim)]
19pub(crate) use self::report::simulated_pk_index_compaction_result;
20pub(crate) use self::report::{
21 IcebergPlanCompletion, IcebergTaskReport, IcebergTaskTracker, PkIndexCompactionResult,
22 ReportSendResult, build_drained_iceberg_task_report, build_iceberg_task_report,
23 build_pk_index_compaction_result, flush_pending_iceberg_task_reports,
24 send_or_buffer_iceberg_task_report,
25};
26use crate::hummock::compactor::iceberg_compaction::iceberg_compactor_runner::IcebergCompactionPlanRunner;
27
28pub(crate) mod iceberg_compactor_runner;
29pub(crate) mod memory;
30pub(crate) mod report;
31
32use std::collections::{HashMap, VecDeque};
33use std::sync::Arc;
34
35use risingwave_pb::id::IcebergCompactionTaskId;
36use tokio::sync::Notify;
37
38use crate::monitor::CompactorMetrics;
39
40pub(crate) type TaskKey = (IcebergCompactionTaskId, usize);
42
43#[derive(Debug, Clone)]
45pub struct IcebergTaskMeta {
46 pub task_id: IcebergCompactionTaskId,
47 pub plan_index: usize,
48 pub required_parallelism: u32,
50 pub memory_reservation_bytes: usize,
53}
54
55#[derive(Debug)]
56pub struct PoppedIcebergTask {
57 pub meta: IcebergTaskMeta,
58 pub runner: Option<IcebergCompactionPlanRunner>,
59}
60
61impl IcebergTaskMeta {
62 fn key(&self) -> TaskKey {
63 (self.task_id, self.plan_index)
64 }
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68struct IcebergTaskResources {
69 required_parallelism: u32,
70 memory_reservation_bytes: usize,
71}
72
73struct IcebergTaskQueueInner {
75 deque: VecDeque<IcebergTaskMeta>,
77 resource_map: HashMap<TaskKey, IcebergTaskResources>,
79 waiting_parallelism_sum: u32,
81 running_parallelism_sum: u32,
83 running_memory_reservation_bytes: usize,
85 runners: HashMap<TaskKey, IcebergCompactionPlanRunner>,
87}
88
89pub struct IcebergTaskQueue {
105 inner: IcebergTaskQueueInner,
106 max_parallelism: u32,
108 pending_parallelism_budget: u32,
110 total_memory_budget_bytes: usize,
112 schedule_notify: Arc<Notify>,
114 metrics: Arc<CompactorMetrics>,
115}
116
117#[derive(Debug, PartialEq, Eq)]
118pub enum PushResult {
119 Added,
120 RejectedCapacity,
122 RejectedTooLarge,
124 RejectedInvalidParallelism,
126 RejectedDuplicate,
128}
129
130impl IcebergTaskQueue {
131 #[cfg(test)]
132 pub fn new(
133 max_parallelism: u32,
134 pending_parallelism_budget: u32,
135 total_memory_budget_bytes: usize,
136 ) -> Self {
137 Self::new_with_metrics(
138 max_parallelism,
139 pending_parallelism_budget,
140 total_memory_budget_bytes,
141 Arc::new(CompactorMetrics::unused()),
142 )
143 }
144
145 pub(crate) fn new_with_metrics(
146 max_parallelism: u32,
147 pending_parallelism_budget: u32,
148 total_memory_budget_bytes: usize,
149 metrics: Arc<CompactorMetrics>,
150 ) -> Self {
151 assert!(max_parallelism > 0, "max_parallelism must be > 0");
152 assert!(
153 pending_parallelism_budget >= max_parallelism,
154 "pending budget should allow at least one task"
155 );
156 assert!(
157 total_memory_budget_bytes > 0,
158 "total memory budget must be > 0"
159 );
160 metrics
161 .iceberg_compaction_memory_budget_bytes
162 .set(i64::try_from(total_memory_budget_bytes).unwrap_or(i64::MAX));
163 metrics
164 .iceberg_compaction_running_memory_reservation_bytes
165 .set(0);
166 Self {
167 inner: IcebergTaskQueueInner {
168 deque: VecDeque::new(),
169 resource_map: HashMap::new(),
170 waiting_parallelism_sum: 0,
171 running_parallelism_sum: 0,
172 running_memory_reservation_bytes: 0,
173 runners: HashMap::new(),
174 },
175 max_parallelism,
176 pending_parallelism_budget,
177 total_memory_budget_bytes,
178 schedule_notify: Arc::new(Notify::new()),
179 metrics,
180 }
181 }
182
183 pub async fn wait_schedulable(&self) -> bool {
188 if self.has_schedulable_tasks() {
190 return true;
191 }
192 self.schedule_notify.notified().await;
194 self.has_schedulable_tasks()
195 }
196
197 fn has_schedulable_tasks(&self) -> bool {
198 if let Some(front_task) = self.inner.deque.front() {
199 let available_parallelism = self
200 .max_parallelism
201 .saturating_sub(self.inner.running_parallelism_sum);
202 let available_memory = self
203 .total_memory_budget_bytes
204 .saturating_sub(self.inner.running_memory_reservation_bytes);
205 available_parallelism >= front_task.required_parallelism
206 && available_memory >= front_task.memory_reservation_bytes
207 } else {
208 false
209 }
210 }
211
212 fn notify_schedulable(&self) {
213 if self.has_schedulable_tasks() {
214 self.schedule_notify.notify_one();
215 }
216 }
217
218 pub fn running_parallelism_sum(&self) -> u32 {
219 self.inner.running_parallelism_sum
220 }
221
222 pub fn waiting_parallelism_sum(&self) -> u32 {
223 self.inner.waiting_parallelism_sum
224 }
225
226 fn available_parallelism(&self) -> u32 {
227 self.max_parallelism
228 .saturating_sub(self.inner.running_parallelism_sum)
229 }
230
231 fn available_memory_reservation_bytes(&self) -> usize {
232 self.total_memory_budget_bytes
233 .saturating_sub(self.inner.running_memory_reservation_bytes)
234 }
235
236 fn update_running_memory_metric(&self) {
237 self.metrics
238 .iceberg_compaction_running_memory_reservation_bytes
239 .set(i64::try_from(self.inner.running_memory_reservation_bytes).unwrap_or(i64::MAX));
240 }
241
242 pub fn push(
246 &mut self,
247 meta: IcebergTaskMeta,
248 runner: Option<IcebergCompactionPlanRunner>,
249 ) -> PushResult {
250 if meta.required_parallelism == 0 {
251 return PushResult::RejectedInvalidParallelism;
252 }
253 if meta.required_parallelism > self.max_parallelism
254 || meta.memory_reservation_bytes > self.total_memory_budget_bytes
255 {
256 return PushResult::RejectedTooLarge;
257 }
258 let key = meta.key();
259
260 if self.inner.resource_map.contains_key(&key) {
262 return PushResult::RejectedDuplicate;
263 }
264
265 let Some(new_parallelism_total) = self
266 .inner
267 .waiting_parallelism_sum
268 .checked_add(meta.required_parallelism)
269 else {
270 return PushResult::RejectedCapacity;
271 };
272 if new_parallelism_total > self.pending_parallelism_budget {
273 return PushResult::RejectedCapacity;
274 }
275 self.inner.resource_map.insert(
276 key,
277 IcebergTaskResources {
278 required_parallelism: meta.required_parallelism,
279 memory_reservation_bytes: meta.memory_reservation_bytes,
280 },
281 );
282 self.inner.waiting_parallelism_sum = new_parallelism_total;
283
284 if let Some(r) = runner {
285 self.inner.runners.insert(key, r);
286 }
287
288 self.inner.deque.push_back(meta);
289
290 self.notify_schedulable();
291 PushResult::Added
292 }
293
294 pub fn pop(&mut self) -> Option<PoppedIcebergTask> {
299 let front = self.inner.deque.front()?;
300 if front.required_parallelism > self.available_parallelism() {
301 return None;
302 }
303 if front.memory_reservation_bytes > self.available_memory_reservation_bytes() {
304 return None;
305 }
306
307 let meta = self.inner.deque.pop_front()?;
308 self.inner.waiting_parallelism_sum = self
309 .inner
310 .waiting_parallelism_sum
311 .saturating_sub(meta.required_parallelism);
312 self.inner.running_parallelism_sum = self
313 .inner
314 .running_parallelism_sum
315 .saturating_add(meta.required_parallelism);
316 self.inner.running_memory_reservation_bytes = self
317 .inner
318 .running_memory_reservation_bytes
319 .checked_add(meta.memory_reservation_bytes)
320 .expect("running memory reservation sum overflowed");
321
322 let key = meta.key();
323 self.update_running_memory_metric();
324 let runner = self.inner.runners.remove(&key);
325 Some(PoppedIcebergTask { meta, runner })
326 }
327
328 pub fn finish_running(&mut self, task_key: TaskKey) -> bool {
332 let Some(resources) = self.inner.resource_map.remove(&task_key) else {
333 tracing::warn!(
334 task_id = %task_key.0,
335 plan_index = task_key.1,
336 "finish_running called for unknown task key, possible bug: double-finish or invalid key"
337 );
338 return false;
339 };
340 self.inner.running_parallelism_sum = self
341 .inner
342 .running_parallelism_sum
343 .saturating_sub(resources.required_parallelism);
344 self.inner.running_memory_reservation_bytes = self
345 .inner
346 .running_memory_reservation_bytes
347 .checked_sub(resources.memory_reservation_bytes)
348 .expect("running memory reservation bookkeeping underflowed");
349 self.inner.runners.remove(&task_key);
350 self.update_running_memory_metric();
351 self.notify_schedulable();
352 true
353 }
354
355 pub fn cancel_waiting_task(&mut self, task_id: IcebergCompactionTaskId) -> usize {
359 let mut retained = VecDeque::with_capacity(self.inner.deque.len());
360 let mut cancelled_parallelism = 0;
361 let mut cancelled_count = 0;
362
363 while let Some(meta) = self.inner.deque.pop_front() {
364 if meta.task_id == task_id {
365 cancelled_parallelism += meta.required_parallelism;
366 cancelled_count += 1;
367 self.inner.resource_map.remove(&meta.key());
368 self.inner.runners.remove(&meta.key());
369 } else {
370 retained.push_back(meta);
371 }
372 }
373
374 self.inner.deque = retained;
375 self.inner.waiting_parallelism_sum = self
376 .inner
377 .waiting_parallelism_sum
378 .saturating_sub(cancelled_parallelism);
379
380 if cancelled_count > 0 {
381 self.notify_schedulable();
382 }
383
384 cancelled_count
385 }
386}
387
388#[cfg(test)]
389mod tests {
390 use super::*;
391
392 fn mk_meta(id: u64, plan_index: usize, p: u32) -> IcebergTaskMeta {
393 IcebergTaskMeta {
394 task_id: id.into(),
395 plan_index,
396 required_parallelism: p,
397 memory_reservation_bytes: 1,
398 }
399 }
400
401 #[test]
402 fn test_basic_push_pop() {
403 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
404 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
405 assert_eq!(q.waiting_parallelism_sum(), 4);
406
407 let popped = q.pop().expect("should pop");
408 assert_eq!(popped.meta.task_id.as_raw_id(), 1);
409 assert_eq!(q.waiting_parallelism_sum(), 0);
410 assert_eq!(q.running_parallelism_sum(), 4);
411
412 assert!(q.finish_running((1.into(), 0)));
413 assert_eq!(q.running_parallelism_sum(), 0);
414 }
415
416 #[test]
417 fn test_fifo_ordering() {
418 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
419 assert_eq!(q.push(mk_meta(1, 0, 2), None), PushResult::Added);
420 assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
421 assert_eq!(q.push(mk_meta(3, 0, 2), None), PushResult::Added);
422
423 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 1);
424 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
425 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 3);
426 }
427
428 #[test]
429 fn test_capacity_reject() {
430 let mut q = IcebergTaskQueue::new(4, 6, usize::MAX);
431 assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
432 assert_eq!(q.push(mk_meta(2, 0, 3), None), PushResult::Added); assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::RejectedCapacity); }
435
436 #[test]
437 fn test_invalid_parallelism() {
438 let mut q = IcebergTaskQueue::new(4, 10, usize::MAX);
439 assert_eq!(
440 q.push(mk_meta(1, 0, 0), None),
441 PushResult::RejectedInvalidParallelism
442 );
443 assert_eq!(q.push(mk_meta(2, 0, 5), None), PushResult::RejectedTooLarge); }
445
446 #[test]
447 fn test_memory_estimate_exceeding_budget_is_rejected() {
448 let mut q = IcebergTaskQueue::new(4, 10, 100);
449 let mut meta = mk_meta(1, 0, 1);
450 meta.memory_reservation_bytes = 101;
451
452 assert_eq!(q.push(meta, None), PushResult::RejectedTooLarge);
453 }
454
455 #[test]
456 fn test_duplicate_key_rejected() {
457 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
458 assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
459 assert_eq!(
461 q.push(mk_meta(1, 0, 5), None),
462 PushResult::RejectedDuplicate
463 );
464 assert_eq!(q.waiting_parallelism_sum(), 3);
466
467 assert_eq!(q.push(mk_meta(1, 1, 2), None), PushResult::Added);
469 assert_eq!(q.waiting_parallelism_sum(), 5);
470
471 let p = q.pop().unwrap();
473 assert_eq!(p.meta.task_id.as_raw_id(), 1);
474 assert_eq!(p.meta.plan_index, 0);
475 q.finish_running((1.into(), 0));
476
477 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
479 }
480
481 #[test]
482 fn test_pop_insufficient_parallelism() {
483 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
484 assert_eq!(q.push(mk_meta(1, 0, 6), None), PushResult::Added);
485 assert_eq!(q.push(mk_meta(2, 0, 4), None), PushResult::Added);
486
487 let p1 = q.pop().unwrap();
488 assert_eq!(p1.meta.task_id.as_raw_id(), 1);
489 assert!(q.pop().is_none());
491
492 assert!(q.finish_running((1.into(), 0)));
494 let p2 = q.pop().unwrap();
495 assert_eq!(p2.meta.task_id.as_raw_id(), 2);
496 }
497
498 #[test]
499 fn test_finish_nonexistent_task() {
500 let mut q = IcebergTaskQueue::new(4, 16, usize::MAX);
501 assert!(!q.finish_running((999.into(), 0)));
502 assert_eq!(q.running_parallelism_sum(), 0);
503 }
504
505 #[test]
506 fn test_double_finish() {
507 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
508 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
509 q.pop().unwrap();
510 assert_eq!(q.running_parallelism_sum(), 4);
511
512 assert!(q.finish_running((1.into(), 0)));
514 assert_eq!(q.running_parallelism_sum(), 0);
515
516 assert!(!q.finish_running((1.into(), 0)));
518 assert_eq!(q.running_parallelism_sum(), 0);
519 }
520
521 #[test]
522 fn test_max_parallelism_boundary() {
523 let mut q = IcebergTaskQueue::new(4, 4, usize::MAX);
525
526 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
528 assert_eq!(q.push(mk_meta(2, 0, 1), None), PushResult::RejectedCapacity);
530
531 let p = q.pop().unwrap();
533 assert_eq!(p.meta.required_parallelism, 4);
534 assert_eq!(q.running_parallelism_sum(), 4);
535
536 assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::Added);
538 assert!(q.pop().is_none());
539 }
540
541 #[test]
542 fn test_same_task_id_multiple_plans() {
543 let mut q = IcebergTaskQueue::new(10, 30, usize::MAX);
544 let task_id = 1u64;
545
546 assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
548 assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
549 assert_eq!(q.push(mk_meta(task_id, 2, 2), None), PushResult::Added);
550 assert_eq!(q.waiting_parallelism_sum(), 9);
551
552 for i in 0..3 {
554 let p = q.pop().unwrap();
555 assert_eq!(p.meta.task_id.as_raw_id(), task_id);
556 assert_eq!(p.meta.plan_index, i);
557 }
558 assert_eq!(q.running_parallelism_sum(), 9);
559
560 assert!(q.finish_running((task_id.into(), 1)));
562 assert_eq!(q.running_parallelism_sum(), 5);
563 assert!(q.finish_running((task_id.into(), 0)));
564 assert!(q.finish_running((task_id.into(), 2)));
565 assert_eq!(q.running_parallelism_sum(), 0);
566 }
567
568 #[test]
569 fn test_cancel_waiting_task_only_removes_waiting_plans() {
570 let mut q = IcebergTaskQueue::new(10, 30, usize::MAX);
571 let task_id = 1u64;
572
573 assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
574 assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
575 assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
576
577 let popped = q.pop().unwrap();
578 assert_eq!(popped.meta.task_id.as_raw_id(), task_id);
579 assert_eq!(popped.meta.plan_index, 0);
580 assert_eq!(q.running_parallelism_sum(), 3);
581 assert_eq!(q.waiting_parallelism_sum(), 6);
582
583 assert_eq!(q.cancel_waiting_task(task_id.into()), 1);
584 assert_eq!(q.running_parallelism_sum(), 3);
585 assert_eq!(q.waiting_parallelism_sum(), 2);
586
587 assert!(q.finish_running((task_id.into(), 0)));
588 let next = q.pop().unwrap();
589 assert_eq!(next.meta.task_id.as_raw_id(), 2);
590 assert_eq!(next.meta.plan_index, 0);
591 }
592
593 #[test]
594 fn test_empty_queue_behavior() {
595 let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
596 assert!(q.pop().is_none());
597 assert!(!q.finish_running((1.into(), 0)));
598 assert_eq!(q.waiting_parallelism_sum(), 0);
599 assert_eq!(q.running_parallelism_sum(), 0);
600 }
601 #[test]
602 fn waiting_memory_does_not_count_against_running_budget() {
603 let mut q = IcebergTaskQueue::new(8, 32, 100);
604 let meta = |id, memory_reservation_bytes| IcebergTaskMeta {
605 task_id: IcebergCompactionTaskId::new(id),
606 plan_index: 0,
607 required_parallelism: 1,
608 memory_reservation_bytes,
609 };
610 assert_eq!(q.push(meta(1, 80), None), PushResult::Added);
611 q.pop().unwrap();
612
613 assert_eq!(q.push(meta(2, 60), None), PushResult::Added);
614 assert!(q.pop().is_none());
615
616 assert!(q.finish_running((1.into(), 0)));
617 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
618 }
619
620 #[test]
621 fn memory_head_of_line_blocking_preserves_fifo() {
622 let mut q = IcebergTaskQueue::new(8, 32, 100);
623 let meta = |id, memory_reservation_bytes| IcebergTaskMeta {
624 task_id: IcebergCompactionTaskId::new(id),
625 plan_index: 0,
626 required_parallelism: 1,
627 memory_reservation_bytes,
628 };
629 assert_eq!(q.push(meta(1, 60), None), PushResult::Added);
630 q.pop().unwrap();
631 assert_eq!(q.push(meta(2, 50), None), PushResult::Added);
632 assert_eq!(q.push(meta(3, 40), None), PushResult::Added);
633
634 assert!(q.pop().is_none());
635 assert!(q.finish_running((1.into(), 0)));
636 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
637 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 3);
638 }
639}