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_iceberg_task_report, build_pk_index_compaction_result,
23 flush_pending_iceberg_task_reports, send_or_buffer_iceberg_task_report,
24};
25use crate::hummock::compactor::iceberg_compaction::iceberg_compactor_runner::IcebergCompactionPlanRunner;
26
27pub(crate) mod iceberg_compactor_runner;
28pub(crate) mod report;
29
30use std::collections::{HashMap, VecDeque};
31use std::sync::Arc;
32
33use risingwave_pb::id::IcebergCompactionTaskId;
34use tokio::sync::Notify;
35
36pub(crate) type TaskKey = (IcebergCompactionTaskId, usize);
38
39#[derive(Debug, Clone)]
41pub struct IcebergTaskMeta {
42 pub task_id: IcebergCompactionTaskId,
43 pub plan_index: usize,
44 pub required_parallelism: u32,
46}
47
48#[derive(Debug)]
49pub struct PoppedIcebergTask {
50 pub meta: IcebergTaskMeta,
51 pub runner: Option<IcebergCompactionPlanRunner>,
52}
53
54impl IcebergTaskMeta {
55 fn key(&self) -> TaskKey {
56 (self.task_id, self.plan_index)
57 }
58}
59
60struct IcebergTaskQueueInner {
62 deque: VecDeque<IcebergTaskMeta>,
64 id_map: HashMap<TaskKey, u32>,
66 waiting_parallelism_sum: u32,
68 running_parallelism_sum: u32,
70 runners: HashMap<TaskKey, IcebergCompactionPlanRunner>,
72}
73
74pub struct IcebergTaskQueue {
88 inner: IcebergTaskQueueInner,
89 max_parallelism: u32,
91 pending_parallelism_budget: u32,
93 schedule_notify: Arc<Notify>,
95}
96
97#[derive(Debug, PartialEq, Eq)]
98pub enum PushResult {
99 Added,
100 RejectedCapacity,
102 RejectedTooLarge,
104 RejectedInvalidParallelism,
106 RejectedDuplicate,
108}
109
110impl IcebergTaskQueue {
111 pub fn new(max_parallelism: u32, pending_parallelism_budget: u32) -> Self {
112 assert!(max_parallelism > 0, "max_parallelism must be > 0");
113 assert!(
114 pending_parallelism_budget >= max_parallelism,
115 "pending budget should allow at least one task"
116 );
117 Self {
118 inner: IcebergTaskQueueInner {
119 deque: VecDeque::new(),
120 id_map: HashMap::new(),
121 waiting_parallelism_sum: 0,
122 running_parallelism_sum: 0,
123 runners: HashMap::new(),
124 },
125 max_parallelism,
126 pending_parallelism_budget,
127 schedule_notify: Arc::new(Notify::new()),
128 }
129 }
130
131 pub async fn wait_schedulable(&self) -> bool {
136 if self.has_schedulable_tasks() {
138 return true;
139 }
140 self.schedule_notify.notified().await;
142 self.has_schedulable_tasks()
143 }
144
145 fn has_schedulable_tasks(&self) -> bool {
146 if let Some(front_task) = self.inner.deque.front() {
147 let available_parallelism = self
148 .max_parallelism
149 .saturating_sub(self.inner.running_parallelism_sum);
150 available_parallelism >= front_task.required_parallelism
151 } else {
152 false
153 }
154 }
155
156 fn notify_schedulable(&self) {
157 if self.has_schedulable_tasks() {
158 self.schedule_notify.notify_one();
159 }
160 }
161
162 pub fn running_parallelism_sum(&self) -> u32 {
163 self.inner.running_parallelism_sum
164 }
165
166 pub fn waiting_parallelism_sum(&self) -> u32 {
167 self.inner.waiting_parallelism_sum
168 }
169
170 fn available_parallelism(&self) -> u32 {
171 self.max_parallelism
172 .saturating_sub(self.inner.running_parallelism_sum)
173 }
174
175 pub fn push(
179 &mut self,
180 meta: IcebergTaskMeta,
181 runner: Option<IcebergCompactionPlanRunner>,
182 ) -> PushResult {
183 if meta.required_parallelism == 0 {
184 return PushResult::RejectedInvalidParallelism;
185 }
186 if meta.required_parallelism > self.max_parallelism {
187 return PushResult::RejectedTooLarge;
188 }
189
190 let key = meta.key();
191
192 if self.inner.id_map.contains_key(&key) {
194 return PushResult::RejectedDuplicate;
195 }
196
197 let new_total = self.inner.waiting_parallelism_sum + meta.required_parallelism;
198 if new_total > self.pending_parallelism_budget {
199 return PushResult::RejectedCapacity;
200 }
201
202 self.inner.id_map.insert(key, meta.required_parallelism);
203 self.inner.waiting_parallelism_sum = new_total;
204
205 if let Some(r) = runner {
206 self.inner.runners.insert(key, r);
207 }
208
209 self.inner.deque.push_back(meta);
210
211 self.notify_schedulable();
212 PushResult::Added
213 }
214
215 pub fn pop(&mut self) -> Option<PoppedIcebergTask> {
220 let front = self.inner.deque.front()?;
221 if front.required_parallelism > self.available_parallelism() {
222 return None;
223 }
224
225 let meta = self.inner.deque.pop_front()?;
226 self.inner.waiting_parallelism_sum = self
227 .inner
228 .waiting_parallelism_sum
229 .saturating_sub(meta.required_parallelism);
230 self.inner.running_parallelism_sum = self
231 .inner
232 .running_parallelism_sum
233 .saturating_add(meta.required_parallelism);
234
235 let runner = self.inner.runners.remove(&meta.key());
236 Some(PoppedIcebergTask { meta, runner })
237 }
238
239 pub fn finish_running(&mut self, task_key: TaskKey) -> bool {
243 let Some(required) = self.inner.id_map.remove(&task_key) else {
244 tracing::warn!(
245 task_id = %task_key.0,
246 plan_index = task_key.1,
247 "finish_running called for unknown task key, possible bug: double-finish or invalid key"
248 );
249 return false;
250 };
251
252 self.inner.running_parallelism_sum =
253 self.inner.running_parallelism_sum.saturating_sub(required);
254 self.inner.runners.remove(&task_key);
255 self.notify_schedulable();
256 true
257 }
258
259 pub fn cancel_waiting_task(&mut self, task_id: IcebergCompactionTaskId) -> usize {
263 let mut retained = VecDeque::with_capacity(self.inner.deque.len());
264 let mut cancelled_parallelism = 0;
265 let mut cancelled_count = 0;
266
267 while let Some(meta) = self.inner.deque.pop_front() {
268 if meta.task_id == task_id {
269 cancelled_parallelism += meta.required_parallelism;
270 cancelled_count += 1;
271 self.inner.id_map.remove(&meta.key());
272 self.inner.runners.remove(&meta.key());
273 } else {
274 retained.push_back(meta);
275 }
276 }
277
278 self.inner.deque = retained;
279 self.inner.waiting_parallelism_sum = self
280 .inner
281 .waiting_parallelism_sum
282 .saturating_sub(cancelled_parallelism);
283
284 if cancelled_count > 0 {
285 self.notify_schedulable();
286 }
287
288 cancelled_count
289 }
290}
291
292#[cfg(test)]
293mod tests {
294 use super::*;
295
296 fn mk_meta(id: u64, plan_index: usize, p: u32) -> IcebergTaskMeta {
297 IcebergTaskMeta {
298 task_id: id.into(),
299 plan_index,
300 required_parallelism: p,
301 }
302 }
303
304 #[test]
305 fn test_basic_push_pop() {
306 let mut q = IcebergTaskQueue::new(8, 32);
307 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
308 assert_eq!(q.waiting_parallelism_sum(), 4);
309
310 let popped = q.pop().expect("should pop");
311 assert_eq!(popped.meta.task_id.as_raw_id(), 1);
312 assert_eq!(q.waiting_parallelism_sum(), 0);
313 assert_eq!(q.running_parallelism_sum(), 4);
314
315 assert!(q.finish_running((1.into(), 0)));
316 assert_eq!(q.running_parallelism_sum(), 0);
317 }
318
319 #[test]
320 fn test_fifo_ordering() {
321 let mut q = IcebergTaskQueue::new(8, 32);
322 assert_eq!(q.push(mk_meta(1, 0, 2), None), PushResult::Added);
323 assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
324 assert_eq!(q.push(mk_meta(3, 0, 2), None), PushResult::Added);
325
326 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 1);
327 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
328 assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 3);
329 }
330
331 #[test]
332 fn test_capacity_reject() {
333 let mut q = IcebergTaskQueue::new(4, 6);
334 assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
335 assert_eq!(q.push(mk_meta(2, 0, 3), None), PushResult::Added); assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::RejectedCapacity); }
338
339 #[test]
340 fn test_invalid_parallelism() {
341 let mut q = IcebergTaskQueue::new(4, 10);
342 assert_eq!(
343 q.push(mk_meta(1, 0, 0), None),
344 PushResult::RejectedInvalidParallelism
345 );
346 assert_eq!(q.push(mk_meta(2, 0, 5), None), PushResult::RejectedTooLarge); }
348
349 #[test]
350 fn test_duplicate_key_rejected() {
351 let mut q = IcebergTaskQueue::new(8, 32);
352 assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
353 assert_eq!(
355 q.push(mk_meta(1, 0, 5), None),
356 PushResult::RejectedDuplicate
357 );
358 assert_eq!(q.waiting_parallelism_sum(), 3);
360
361 assert_eq!(q.push(mk_meta(1, 1, 2), None), PushResult::Added);
363 assert_eq!(q.waiting_parallelism_sum(), 5);
364
365 let p = q.pop().unwrap();
367 assert_eq!(p.meta.task_id.as_raw_id(), 1);
368 assert_eq!(p.meta.plan_index, 0);
369 q.finish_running((1.into(), 0));
370
371 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
373 }
374
375 #[test]
376 fn test_pop_insufficient_parallelism() {
377 let mut q = IcebergTaskQueue::new(8, 32);
378 assert_eq!(q.push(mk_meta(1, 0, 6), None), PushResult::Added);
379 assert_eq!(q.push(mk_meta(2, 0, 4), None), PushResult::Added);
380
381 let p1 = q.pop().unwrap();
382 assert_eq!(p1.meta.task_id.as_raw_id(), 1);
383 assert!(q.pop().is_none());
385
386 assert!(q.finish_running((1.into(), 0)));
388 let p2 = q.pop().unwrap();
389 assert_eq!(p2.meta.task_id.as_raw_id(), 2);
390 }
391
392 #[test]
393 fn test_finish_nonexistent_task() {
394 let mut q = IcebergTaskQueue::new(4, 16);
395 assert!(!q.finish_running((999.into(), 0)));
396 assert_eq!(q.running_parallelism_sum(), 0);
397 }
398
399 #[test]
400 fn test_double_finish() {
401 let mut q = IcebergTaskQueue::new(8, 32);
402 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
403 q.pop().unwrap();
404 assert_eq!(q.running_parallelism_sum(), 4);
405
406 assert!(q.finish_running((1.into(), 0)));
408 assert_eq!(q.running_parallelism_sum(), 0);
409
410 assert!(!q.finish_running((1.into(), 0)));
412 assert_eq!(q.running_parallelism_sum(), 0);
413 }
414
415 #[test]
416 fn test_max_parallelism_boundary() {
417 let mut q = IcebergTaskQueue::new(4, 4);
419
420 assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
422 assert_eq!(q.push(mk_meta(2, 0, 1), None), PushResult::RejectedCapacity);
424
425 let p = q.pop().unwrap();
427 assert_eq!(p.meta.required_parallelism, 4);
428 assert_eq!(q.running_parallelism_sum(), 4);
429
430 assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::Added);
432 assert!(q.pop().is_none());
433 }
434
435 #[test]
436 fn test_same_task_id_multiple_plans() {
437 let mut q = IcebergTaskQueue::new(10, 30);
438 let task_id = 1u64;
439
440 assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
442 assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
443 assert_eq!(q.push(mk_meta(task_id, 2, 2), None), PushResult::Added);
444 assert_eq!(q.waiting_parallelism_sum(), 9);
445
446 for i in 0..3 {
448 let p = q.pop().unwrap();
449 assert_eq!(p.meta.task_id.as_raw_id(), task_id);
450 assert_eq!(p.meta.plan_index, i);
451 }
452 assert_eq!(q.running_parallelism_sum(), 9);
453
454 assert!(q.finish_running((task_id.into(), 1)));
456 assert_eq!(q.running_parallelism_sum(), 5);
457 assert!(q.finish_running((task_id.into(), 0)));
458 assert!(q.finish_running((task_id.into(), 2)));
459 assert_eq!(q.running_parallelism_sum(), 0);
460 }
461
462 #[test]
463 fn test_cancel_waiting_task_only_removes_waiting_plans() {
464 let mut q = IcebergTaskQueue::new(10, 30);
465 let task_id = 1u64;
466
467 assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
468 assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
469 assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
470
471 let popped = q.pop().unwrap();
472 assert_eq!(popped.meta.task_id.as_raw_id(), task_id);
473 assert_eq!(popped.meta.plan_index, 0);
474 assert_eq!(q.running_parallelism_sum(), 3);
475 assert_eq!(q.waiting_parallelism_sum(), 6);
476
477 assert_eq!(q.cancel_waiting_task(task_id.into()), 1);
478 assert_eq!(q.running_parallelism_sum(), 3);
479 assert_eq!(q.waiting_parallelism_sum(), 2);
480
481 assert!(q.finish_running((task_id.into(), 0)));
482 let next = q.pop().unwrap();
483 assert_eq!(next.meta.task_id.as_raw_id(), 2);
484 assert_eq!(next.meta.plan_index, 0);
485 }
486
487 #[test]
488 fn test_empty_queue_behavior() {
489 let mut q = IcebergTaskQueue::new(8, 32);
490 assert!(q.pop().is_none());
491 assert!(!q.finish_running((1.into(), 0)));
492 assert_eq!(q.waiting_parallelism_sum(), 0);
493 assert_eq!(q.running_parallelism_sum(), 0);
494 }
495}