Skip to main content

servers/batcher/table/
pending_worker.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
40/// The sender handle does not own the worker's processing state.
41pub(in crate::batcher::table) struct PendingWorker {
42    pub tx: mpsc::Sender<WorkerCommand>,
43}
44
45/// Identifies submissions that can share a batch and overlapping writes.
46#[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        // Handles are retained only for the table identity and schema transition barrier. Dropping
80        // them on shutdown detaches writes, matching the Prom worker contract.
81        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                                // Finish writes from the previous table definition before switching.
99                                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                        // Keep table identity while writes are in flight. Closing
122                        // still allows previously reserved sends to be received.
123                        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}