Skip to main content

servers/batcher/
logical_table.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
15mod batch;
16mod batch_convert;
17mod flow_notifier;
18mod pending_worker;
19mod region_write;
20mod tables;
21#[cfg(test)]
22mod test_util;
23
24use std::collections::{HashMap, HashSet};
25use std::future::{Future, ready};
26use std::num::NonZeroUsize;
27use std::sync::Arc;
28use std::time::Duration;
29
30use api::v1::RowInsertRequests;
31use catalog::CatalogManagerRef;
32use common_batcher::flush_limiter::FlushLimiter;
33use common_batcher::flush_policy::timing::TimingFlushPolicy;
34use common_batcher::request_limiter::RequestLimiter;
35use common_batcher::worker_registry::WorkerRegistry;
36use common_meta::cache::TableFlownodeSetCacheRef;
37use common_meta::node_manager::NodeManagerRef;
38use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
39use common_time::timestamp::TimeUnit;
40use meter_core::data::MeterRecord;
41use meter_macros::write_meter;
42use partition::manager::PartitionRuleManagerRef;
43use session::context::QueryContextRef;
44use snafu::ResultExt;
45use store_api::metric_engine_consts::{LOGICAL_TABLE_METADATA_KEY, METRIC_ENGINE_NAME};
46use tokio::sync::{Semaphore, broadcast, mpsc, oneshot};
47
48use crate::batcher::flow_notifier::{FlowNotifier, start_flow_notification_worker};
49pub use crate::batcher::logical_table::batch::flush_batch_physical;
50pub use crate::batcher::logical_table::batch_convert::{RecordBatchWithTsIdx, TableBatch};
51use crate::batcher::logical_table::pending_worker::{
52    PendingWorker, WorkerCommand, remove_worker_if_same_channel, start_worker,
53};
54pub use crate::batcher::logical_table::region_write::{
55    PhysicalFlushCatalogProvider, PhysicalFlushNodeRequester, PhysicalFlushPartitionProvider,
56    PhysicalTableMetadata,
57};
58pub use crate::batcher::logical_table::tables::{
59    PendingRowsSchemaAlterer, PendingRowsSchemaAltererRef,
60};
61use crate::batcher::pending_rows_batch_sync_enabled;
62use crate::error;
63use crate::error::{Error, Result};
64use crate::metrics::{
65    FLOW_NOTIFICATION_DROPPED, PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED, PENDING_WORKERS,
66};
67
68const PHYSICAL_TABLE_KEY: &str = "physical_table";
69
70const WORKER_IDLE_TIMEOUT_MULTIPLIER: u32 = 3;
71
72#[derive(Debug, Clone, Hash, Eq, PartialEq)]
73pub(crate) struct BatchKey {
74    pub(crate) catalog: String,
75    pub(crate) schema: String,
76    pub(crate) physical_table: String,
77    pub(crate) skip_wal: bool,
78}
79
80// Requests can share a batch only when their write target and WAL policy match.
81pub(crate) fn batch_key_from_ctx(ctx: &QueryContextRef) -> BatchKey {
82    let physical_table = ctx
83        .extension(PHYSICAL_TABLE_KEY)
84        .unwrap_or(GREPTIME_PHYSICAL_TABLE)
85        .to_string();
86    BatchKey {
87        catalog: ctx.current_catalog().to_string(),
88        schema: ctx.current_schema(),
89        physical_table,
90        skip_wal: ctx.skip_wal(),
91    }
92}
93
94/// Prometheus remote write pending rows batcher.
95pub struct LogicalTablePendingRowsBatcher {
96    workers: Arc<WorkerRegistry<BatchKey, WorkerCommand>>,
97    flush_interval: Duration,
98    flush_policy: TimingFlushPolicy,
99    partition_manager: PartitionRuleManagerRef,
100    node_manager: NodeManagerRef,
101    catalog_manager: CatalogManagerRef,
102    flow_notification_tx: FlowNotifier,
103    flush_limiter: FlushLimiter,
104    request_limiter: RequestLimiter,
105    worker_channel_capacity: usize,
106    prom_store_with_metric_engine: bool,
107    schema_alterer: PendingRowsSchemaAltererRef,
108    pending_rows_batch_sync: bool,
109    shutdown: broadcast::Sender<()>,
110}
111
112impl LogicalTablePendingRowsBatcher {
113    #[allow(clippy::too_many_arguments)]
114    pub fn try_new(
115        partition_manager: PartitionRuleManagerRef,
116        node_manager: NodeManagerRef,
117        catalog_manager: CatalogManagerRef,
118        table_flownode_set_cache: TableFlownodeSetCacheRef,
119        prom_store_with_metric_engine: bool,
120        schema_alterer: PendingRowsSchemaAltererRef,
121        flush_interval: Duration,
122        max_batch_rows: usize,
123        max_concurrent_flushes: usize,
124        worker_channel_capacity: usize,
125        max_inflight_requests: usize,
126        flow_notification_queue_capacity: NonZeroUsize,
127    ) -> Option<Arc<Self>> {
128        if worker_channel_capacity == 0 || worker_channel_capacity > Semaphore::MAX_PERMITS {
129            return None;
130        }
131
132        let flush_policy = TimingFlushPolicy::try_new(flush_interval, max_batch_rows)?;
133        let flush_limiter = FlushLimiter::try_new(max_concurrent_flushes)?;
134
135        let request_limiter = RequestLimiter::try_new(max_inflight_requests)?;
136        let (flow_notification_tx, flow_notification_rx) = FlowNotifier::try_new(
137            flow_notification_queue_capacity.get(),
138            FLOW_NOTIFICATION_DROPPED.clone(),
139        )?;
140
141        let (shutdown, _) = broadcast::channel(1);
142        let pending_rows_batch_sync = pending_rows_batch_sync_enabled();
143        let workers = Arc::new(WorkerRegistry::new());
144        PENDING_WORKERS.set(0);
145        start_flow_notification_worker(
146            flow_notification_rx,
147            table_flownode_set_cache,
148            node_manager.clone(),
149        );
150
151        Some(Arc::new(Self {
152            workers,
153            flush_interval,
154            flush_policy,
155            partition_manager,
156            node_manager,
157            catalog_manager,
158            flow_notification_tx,
159            prom_store_with_metric_engine,
160            schema_alterer,
161            flush_limiter,
162            request_limiter,
163            worker_channel_capacity,
164            pending_rows_batch_sync,
165            shutdown,
166        }))
167    }
168}
169
170impl LogicalTablePendingRowsBatcher {
171    pub async fn submit(&self, requests: RowInsertRequests, ctx: QueryContextRef) -> Result<u64> {
172        self.submit_with(requests, ctx, |_| ready(Ok(())))
173            .await
174            .map(|(rows, ())| rows)
175    }
176
177    /// Returns the physical metric table's time index unit resolved from
178    /// `ctx`, defaulting to millisecond when the table does not exist yet
179    /// (the schema alterer auto-creates it as millisecond).
180    async fn physical_time_index_unit_or_default(&self, ctx: &QueryContextRef) -> TimeUnit {
181        let key = batch_key_from_ctx(ctx);
182        let Ok(Some(table)) = self
183            .catalog_manager
184            .table(&key.catalog, &key.schema, &key.physical_table, None)
185            .await
186        else {
187            return TimeUnit::Millisecond;
188        };
189        table
190            .table_info()
191            .meta
192            .schema
193            .timestamp_column()
194            .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit()))
195            .unwrap_or(TimeUnit::Millisecond)
196    }
197
198    /// Returns whether the bulk path can accept `batches`: every existing
199    /// destination table must be a metric logical table bound to the
200    /// physical table selected by its context. A destination bound to
201    /// another physical table would be flushed through the selected
202    /// physical's regions, silently misplacing its rows, so such requests
203    /// must stay on the ordinary insert path (which routes per destination).
204    /// New tables are always fine: they are created on the selected physical
205    /// table. Time index units need no check here — the bulk encode converts
206    /// each request to its destination's unit. Destinations are resolved
207    /// once per distinct (schema, table).
208    pub(crate) async fn accepts_bulk_destinations(
209        &self,
210        batches: impl Iterator<Item = &(QueryContextRef, RowInsertRequests)>,
211    ) -> bool {
212        // One request can select different physical tables per batch (e.g.
213        // per-series physical-table labels), so the dedupe key includes the
214        // selected physical table: every distinct (schema, table, physical)
215        // triple is validated against the table's actual binding.
216        let mut checked = HashSet::new();
217        // For missing tables, one request must not select two different
218        // physical tables: the batcher would create the table through one
219        // selection and flush its rows through the other's regions.
220        let mut missing_selections: HashMap<(String, String), String> = HashMap::new();
221        for (ctx, requests) in batches {
222            let physical_table = batch_key_from_ctx(ctx).physical_table;
223            let schema = ctx.current_schema();
224            for request in &requests.inserts {
225                if !checked.insert((
226                    schema.clone(),
227                    request.table_name.clone(),
228                    physical_table.clone(),
229                )) {
230                    continue;
231                }
232                let Ok(Some(table)) = self
233                    .catalog_manager
234                    .table(ctx.current_catalog(), &schema, &request.table_name, None)
235                    .await
236                else {
237                    // New table: created on the selected physical table, but
238                    // a conflicting selection within the same request cannot
239                    // be batched.
240                    if missing_selections
241                        .insert(
242                            (schema.clone(), request.table_name.clone()),
243                            physical_table.clone(),
244                        )
245                        .is_some_and(|previous| previous != physical_table)
246                    {
247                        return false;
248                    }
249                    continue;
250                };
251                let info = table.table_info();
252                if info.meta.engine != METRIC_ENGINE_NAME
253                    || info
254                        .meta
255                        .options
256                        .extra_options
257                        .get(LOGICAL_TABLE_METADATA_KEY)
258                        .map(String::as_str)
259                        != Some(physical_table.as_str())
260                {
261                    return false;
262                }
263            }
264        }
265        true
266    }
267
268    /// Submits with request-level accounting after schema preparation and before
269    /// queue admission. Acknowledgement follows the global batching policy.
270    pub async fn submit_with<T, F>(
271        &self,
272        requests: RowInsertRequests,
273        ctx: QueryContextRef,
274        after_prepare: impl FnOnce(RowInsertRequests) -> F,
275    ) -> Result<(u64, T)>
276    where
277        F: Future<Output = Result<T>> + Send,
278    {
279        let (table_batches, total_rows) = {
280            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
281                .with_label_values(&["submit_build_and_align"])
282                .start_timer();
283            self.build_and_align_table_batches(&requests, &ctx).await?
284        };
285        let prepared = after_prepare(requests).await?;
286        if total_rows == 0 {
287            return Ok((0, prepared));
288        }
289
290        // Flushes dispatch directly to datanodes, so admit once before enqueueing.
291        write_meter!(MeterRecord::new(
292            ctx.current_catalog().to_string(),
293            ctx.current_schema(),
294            0,
295            ctx.write_rows_to_admit(
296                ctx.current_catalog(),
297                &ctx.current_schema(),
298                total_rows as u64
299            ),
300            ctx.channel() as u8,
301        ))
302        .await
303        .context(error::WriteRejectedSnafu)?;
304
305        let permit = {
306            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
307                .with_label_values(&["submit_acquire_inflight_permit"])
308                .start_timer();
309            self.request_limiter
310                .acquire()
311                .await
312                .map_err(|_| error::BatcherChannelClosedSnafu.build())?
313        };
314
315        let (response_tx, response_rx) = oneshot::channel();
316
317        let batch_key = batch_key_from_ctx(&ctx);
318        let mut cmd = Some(WorkerCommand::Submit {
319            table_batches,
320            total_rows,
321            ctx,
322            response_tx,
323            _permit: permit,
324        });
325
326        {
327            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
328                .with_label_values(&["submit_send_to_worker"])
329                .start_timer();
330
331            for _ in 0..2 {
332                let worker = self.get_or_spawn_worker(batch_key.clone()).await;
333                let Some(worker_cmd) = cmd.take() else {
334                    break;
335                };
336
337                match worker.tx.send(worker_cmd).await {
338                    Ok(()) => break,
339                    Err(err) => {
340                        cmd = Some(err.0);
341                        remove_worker_if_same_channel(
342                            self.workers.as_ref(),
343                            &batch_key,
344                            &worker.tx,
345                        )
346                        .await;
347                    }
348                }
349            }
350
351            if cmd.is_some() {
352                return Err(Error::BatcherChannelClosed);
353            }
354        }
355
356        if self.pending_rows_batch_sync {
357            let result = {
358                let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
359                    .with_label_values(&["submit_wait_flush_result"])
360                    .start_timer();
361                response_rx
362                    .await
363                    .map_err(|_| error::BatcherChannelClosedSnafu.build())?
364            };
365            result
366                .context(error::SubmitBatchSnafu)
367                .map(|()| (total_rows as u64, prepared))
368        } else {
369            Ok((total_rows as u64, prepared))
370        }
371    }
372}
373
374impl LogicalTablePendingRowsBatcher {
375    async fn get_or_spawn_worker(&self, key: BatchKey) -> PendingWorker {
376        let (tx, receiver) = self
377            .workers
378            .get_or_create(key.clone(), self.worker_channel_capacity)
379            .await;
380        if let Some(rx) = receiver {
381            self.spawn_worker(key, tx.clone(), rx);
382            PENDING_WORKERS.set(self.workers.len().await as i64);
383        }
384        PendingWorker { tx }
385    }
386}
387
388impl LogicalTablePendingRowsBatcher {
389    fn spawn_worker(
390        &self,
391        key: BatchKey,
392        tx: mpsc::Sender<WorkerCommand>,
393        rx: mpsc::Receiver<WorkerCommand>,
394    ) {
395        let worker_idle_timeout = self
396            .flush_interval
397            .checked_mul(WORKER_IDLE_TIMEOUT_MULTIPLIER)
398            .unwrap_or(self.flush_interval);
399
400        start_worker(
401            key,
402            tx,
403            self.workers.clone(),
404            rx,
405            self.shutdown.clone(),
406            self.partition_manager.clone(),
407            self.node_manager.clone(),
408            self.catalog_manager.clone(),
409            self.flow_notification_tx.clone(),
410            worker_idle_timeout,
411            self.flush_policy,
412            self.flush_limiter.clone(),
413        );
414    }
415}
416
417impl Drop for LogicalTablePendingRowsBatcher {
418    fn drop(&mut self) {
419        let _ = self.shutdown.send(());
420    }
421}
422
423#[cfg(test)]
424mod tests {
425    use std::sync::Arc;
426
427    use arrow::array::{StringArray, TimestampMillisecondArray};
428    use arrow::datatypes::{DataType as ArrowDataType, Field, Schema as ArrowSchema};
429    use arrow::record_batch::RecordBatch;
430    use common_query::prelude::greptime_timestamp;
431
432    use crate::batcher::logical_table::batch_convert::{RecordBatchWithTsIdx, TableBatch};
433    use crate::batcher::logical_table::flow_notifier::extract_timestamps;
434    use crate::batcher::logical_table::test_util::{mock_aligned_tag_batch, mock_tag_batch};
435
436    #[test]
437    fn test_extract_timestamps_appends_non_null_batches_in_order() {
438        let table_batch = TableBatch {
439            table_name: "cpu".to_string(),
440            table_id: 42,
441            batches: vec![
442                mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0),
443                mock_aligned_tag_batch("tag1", "host-1", 2000, 2.0),
444            ],
445            row_count: 2,
446        };
447
448        assert_eq!(vec![1000, 2000], extract_timestamps(&table_batch));
449    }
450
451    #[test]
452    fn test_extract_timestamps_omits_nulls_and_retains_order() {
453        let table_batch = TableBatch {
454            table_name: "cpu".to_string(),
455            table_id: 42,
456            batches: vec![
457                mock_timestamp_batch(vec![Some(1000), None, Some(3000)]),
458                mock_timestamp_batch(vec![None, Some(5000)]),
459            ],
460            row_count: 5,
461        };
462
463        assert_eq!(vec![1000, 3000, 5000], extract_timestamps(&table_batch));
464    }
465
466    #[test]
467    fn test_record_batch_with_ts_idx_rejects_out_of_bounds_index() {
468        let batch = mock_tag_batch("tag1", "host-1", 1000, 1.0);
469
470        assert!(RecordBatchWithTsIdx::try_new(batch, 3).is_err());
471    }
472
473    #[test]
474    fn test_record_batch_with_ts_idx_rejects_non_timestamp_column() {
475        let batch = mock_tag_batch("tag1", "host-1", 1000, 1.0);
476
477        assert!(RecordBatchWithTsIdx::try_new(batch, 1).is_err());
478    }
479
480    #[test]
481    fn test_extract_timestamps_supports_per_batch_timestamp_indices() {
482        let timestamp_first = RecordBatch::try_new(
483            Arc::new(ArrowSchema::new(vec![
484                Field::new(
485                    "ts",
486                    ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
487                    false,
488                ),
489                Field::new("host", ArrowDataType::Utf8, true),
490            ])),
491            vec![
492                Arc::new(TimestampMillisecondArray::from(vec![1000, 2000])),
493                Arc::new(StringArray::from(vec!["host-1", "host-2"])),
494            ],
495        )
496        .unwrap();
497        let timestamp_second = RecordBatch::try_new(
498            Arc::new(ArrowSchema::new(vec![
499                Field::new("host", ArrowDataType::Utf8, true),
500                Field::new(
501                    "timestamp",
502                    ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
503                    false,
504                ),
505            ])),
506            vec![
507                Arc::new(StringArray::from(vec!["host-3", "host-4"])),
508                Arc::new(TimestampMillisecondArray::from(vec![3000, 4000])),
509            ],
510        )
511        .unwrap();
512        let table_batch = TableBatch {
513            table_name: "cpu".to_string(),
514            table_id: 42,
515            batches: vec![
516                RecordBatchWithTsIdx::try_new(timestamp_first, 0).unwrap(),
517                RecordBatchWithTsIdx::try_new(timestamp_second, 1).unwrap(),
518            ],
519            row_count: 4,
520        };
521
522        assert_eq!(
523            vec![1000, 2000, 3000, 4000],
524            extract_timestamps(&table_batch)
525        );
526    }
527
528    fn mock_timestamp_batch(timestamps: Vec<Option<i64>>) -> RecordBatchWithTsIdx {
529        let batch = RecordBatch::try_new(
530            Arc::new(ArrowSchema::new(vec![Field::new(
531                greptime_timestamp(),
532                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
533                true,
534            )])),
535            vec![Arc::new(TimestampMillisecondArray::from(timestamps))],
536        )
537        .unwrap();
538        RecordBatchWithTsIdx::try_new(batch, 0).unwrap()
539    }
540    #[test]
541    fn test_batch_key_groups_by_skip_wal() {
542        use std::collections::HashMap;
543
544        use crate::batcher::logical_table::batch_key_from_ctx;
545
546        let wal_ctx = session::context::QueryContext::arc();
547        let skip_wal_ctx = session::context::QueryContext::arc();
548        skip_wal_ctx.set_skip_wal(true);
549        let another_skip_wal_ctx = session::context::QueryContext::arc();
550        another_skip_wal_ctx.set_skip_wal(true);
551
552        let mut batches = HashMap::new();
553        *batches.entry(batch_key_from_ctx(&wal_ctx)).or_insert(0) += 1;
554        *batches
555            .entry(batch_key_from_ctx(&skip_wal_ctx))
556            .or_insert(0) += 1;
557        *batches
558            .entry(batch_key_from_ctx(&another_skip_wal_ctx))
559            .or_insert(0) += 1;
560        *batches
561            .entry(batch_key_from_ctx(&session::context::QueryContext::arc()))
562            .or_insert(0) += 1;
563
564        assert_eq!(batches.len(), 2);
565        assert_eq!(batches[&batch_key_from_ctx(&wal_ctx)], 2);
566        assert_eq!(batches[&batch_key_from_ctx(&skip_wal_ctx)], 2);
567    }
568}