servers/batcher/table/
pending_worker.rs1use std::sync::Arc;
16use std::time::Duration;
17
18use arrow::datatypes::SchemaRef;
19use common_batcher::flush_limiter::FlushLimiter;
20use common_batcher::flush_policy::FlushTrigger;
21use common_batcher::flush_policy::timing::TimingFlushPolicy;
22use common_batcher::pending_worker::PendingWorker as PendingCore;
23use common_runtime::spawn_global;
24use common_telemetry::warn;
25use operator::insert::Inserter;
26use tokio::sync::{broadcast, mpsc};
27use tokio::task::JoinHandle;
28use tokio::time::{Instant, sleep_until};
29
30use crate::batcher::table::batch::{Batch, flush_batch};
31use crate::batcher::table::flow_notifier::FlowNotifier;
32use crate::batcher::table::metrics::{PENDING_BATCHES, PENDING_ROWS, PENDING_WORKERS};
33use crate::batcher::table::pending_batch::PendingBatch;
34use crate::batcher::table::{BatchKey, PendingWorkers};
35
36pub(in crate::batcher::table) enum WorkerCommand {
37 Submit(PendingBatch),
38}
39
40pub(in crate::batcher::table) struct PendingWorker {
42 pub tx: mpsc::Sender<WorkerCommand>,
43}
44
45#[derive(PartialEq)]
47struct BatchIdentity {
48 table_id: u32,
49 table_version: u64,
50 schema: SchemaRef,
51}
52
53impl BatchIdentity {
54 fn new(batch: &PendingBatch) -> Self {
55 Self {
56 table_id: batch.table_info.table_id(),
57 table_version: batch.table_info.ident.version,
58 schema: batch.batch.schema(),
59 }
60 }
61}
62
63#[allow(clippy::too_many_arguments)]
64pub(in crate::batcher::table) fn start_worker(
65 key: BatchKey,
66 worker_tx: mpsc::Sender<WorkerCommand>,
67 workers: Arc<PendingWorkers>,
68 mut rx: mpsc::Receiver<WorkerCommand>,
69 shutdown: broadcast::Sender<()>,
70 flush_policy: TimingFlushPolicy,
71 flush_limiter: FlushLimiter,
72 inserter: Arc<Inserter>,
73 notifier: FlowNotifier,
74 idle_timeout: Duration,
75) {
76 spawn_global(async move {
77 let mut pending = PendingCore::new(flush_policy);
78 let mut identity: Option<BatchIdentity> = None;
79 let mut flushes: Vec<JoinHandle<()>> = Vec::new();
82 let mut shutdown_rx = shutdown.subscribe();
83 let mut closing = false;
84 let idle_timer = sleep_until(Instant::now() + idle_timeout);
85 tokio::pin!(idle_timer);
86 loop {
87 let batch = tokio::select! {
88 command = rx.recv() => {
89 match command {
90 Some(WorkerCommand::Submit(batch)) => {
91 idle_timer.as_mut().reset(Instant::now() + idle_timeout);
92 let next_identity = BatchIdentity::new(&batch);
93 if identity.as_ref().is_some_and(|identity| *identity != next_identity) {
94 if let Some(old) = drain_batch(&mut pending, None)
95 && let Some(task) = spawn_flush(old, &flush_limiter, inserter.clone(), notifier.clone()).await {
96 flushes.push(task);
97 }
98 for task in flushes.drain(..) {
100 if let Err(error) = task.await { warn!(error; "Failed to join old-schema batch flush"); }
101 }
102 }
103 identity = Some(next_identity);
104 let rows = batch.batch.num_rows();
105 if pending.is_empty() { PENDING_BATCHES.inc(); }
106 pending.submit(batch, rows);
107 PENDING_ROWS.add(rows as i64);
108 drain_batch(&mut pending, Some(FlushTrigger::Submission))
109 }
110 None => {
111 if let Some(batch) = drain_batch(&mut pending, None) {
112 flush_batch(batch, inserter.clone(), notifier.clone()).await;
113 }
114 break;
115 }
116 }
117 }
118 _ = pending.wait_flush() => drain_batch(&mut pending, Some(FlushTrigger::Deadline)),
119 _ = &mut idle_timer, if !closing => {
120 if pending.is_empty() && rx.is_empty() && flushes.iter().all(JoinHandle::is_finished) {
121 rx.close();
124 closing = true;
125 }
126 idle_timer.as_mut().reset(Instant::now() + idle_timeout);
127 None
128 }
129 _ = shutdown_rx.recv() => {
130 if let Some(batch) = drain_batch(&mut pending, None) {
131 flush_batch(batch, inserter.clone(), notifier.clone()).await;
132 }
133 break;
134 }
135 };
136 if let Some(batch) = batch {
137 flushes.retain(|task| !task.is_finished());
138 if let Some(task) =
139 spawn_flush(batch, &flush_limiter, inserter.clone(), notifier.clone()).await
140 {
141 flushes.push(task);
142 }
143 }
144 }
145 if workers.remove_if_same(&key, &worker_tx).await {
146 PENDING_WORKERS.set(workers.len().await as i64);
147 }
148 });
149}
150
151fn drain_batch(
152 pending: &mut PendingCore<PendingBatch, TimingFlushPolicy>,
153 trigger: Option<FlushTrigger>,
154) -> Option<Batch> {
155 let total_rows = pending.total_rows();
156 let submissions = match trigger {
157 Some(trigger) => pending.take_ready(trigger)?,
158 None => pending.take_pending()?,
159 };
160 PENDING_ROWS.sub(total_rows as i64);
161 PENDING_BATCHES.dec();
162 Some(Batch {
163 submissions,
164 total_rows,
165 })
166}
167
168async fn spawn_flush(
169 batch: Batch,
170 limiter: &FlushLimiter,
171 inserter: Arc<Inserter>,
172 notifier: FlowNotifier,
173) -> Option<JoinHandle<()>> {
174 match limiter.acquire().await {
175 Ok(permit) => Some(spawn_global(async move {
176 let _permit = permit;
177 flush_batch(batch, inserter, notifier).await;
178 })),
179 Err(error) => {
180 warn!(error; "Flush limiter closed, flushing inline");
181 flush_batch(batch, inserter, notifier).await;
182 None
183 }
184 }
185}