Skip to main content

servers/batcher/table/
batch.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::Instant;
17
18use arrow::compute::concat_batches;
19use arrow::record_batch::RecordBatch;
20use common_telemetry::error;
21use operator::error::{ComputeArrowSnafu, Result, UnexpectedSnafu};
22use operator::insert::Inserter;
23use operator::metrics::DIST_INGEST_ROW_COUNT;
24use snafu::ResultExt;
25use table::metadata::TableInfoRef;
26
27use crate::batcher::table::flow_notifier::FlowNotifier;
28use crate::batcher::table::metrics::{
29    FLUSH_DROPPED_ROWS, FLUSH_ELAPSED, FLUSH_FAILURES, FLUSH_ROWS, FLUSH_TOTAL,
30};
31use crate::batcher::table::pending_batch::{PendingBatch, notify_batches};
32
33/// A detached batch owns every submission until its completion is reported.
34pub(in crate::batcher::table) struct Batch {
35    pub submissions: Vec<PendingBatch>,
36    pub total_rows: usize,
37}
38
39pub(in crate::batcher::table) async fn flush_batch(
40    batch: Batch,
41    inserter: Arc<Inserter>,
42    notifier: FlowNotifier,
43) {
44    let started = Instant::now();
45    let result = send_batch(&batch, &inserter).await;
46    FLUSH_ELAPSED.observe(started.elapsed().as_secs_f64());
47    match result {
48        Ok((table, combined)) => {
49            FLUSH_TOTAL.inc();
50            FLUSH_ROWS.observe(batch.total_rows as f64);
51            DIST_INGEST_ROW_COUNT
52                .with_label_values(&[batch.submissions[0].ctx.get_db_string().as_str()])
53                .inc_by(batch.total_rows as u64);
54            notify_batches(batch.submissions, Ok(()));
55            // Admission is best-effort and only occurs after affected-row validation.
56            notifier.notify(table, &combined);
57        }
58        Err(error) => {
59            error!(error; "Failed to flush table batch, rows: {}", batch.total_rows);
60            FLUSH_FAILURES.inc();
61            FLUSH_DROPPED_ROWS.inc_by(batch.total_rows as u64);
62            notify_batches(batch.submissions, Err(Arc::new(error)));
63        }
64    }
65}
66
67async fn send_batch(batch: &Batch, inserter: &Inserter) -> Result<(TableInfoRef, RecordBatch)> {
68    let first = batch.submissions.first().ok_or_else(|| {
69        UnexpectedSnafu {
70            violated: "cannot execute an empty flush".to_string(),
71        }
72        .build()
73    })?;
74    let combined = combine_batches(batch.submissions.iter().map(|pending| &pending.batch))?;
75    let affected_rows = inserter
76        .flush_bulk_batch(
77            first.table_info.clone(),
78            combined.clone(),
79            first.ctx.clone(),
80        )
81        .await?;
82    validate_affected_rows(batch.total_rows, affected_rows)?;
83    Ok((first.table_info.clone(), combined))
84}
85
86fn combine_batches<'a>(batches: impl IntoIterator<Item = &'a RecordBatch>) -> Result<RecordBatch> {
87    let batches = batches.into_iter().collect::<Vec<_>>();
88    let Some(first) = batches.first() else {
89        return UnexpectedSnafu {
90            violated: "cannot combine an empty flush".to_string(),
91        }
92        .fail();
93    };
94    let schema = first.schema();
95    if batches.iter().any(|batch| batch.schema() != schema) {
96        return UnexpectedSnafu {
97            violated: "cannot combine different Arrow schemas".to_string(),
98        }
99        .fail();
100    }
101    concat_batches(&schema, batches).context(ComputeArrowSnafu)
102}
103
104fn validate_affected_rows(expected: usize, actual: usize) -> Result<()> {
105    if expected != actual {
106        return UnexpectedSnafu {
107            violated: format!("batched write affected {actual} rows, expected {expected}; individual results cannot be attributed"),
108        }.fail();
109    }
110    Ok(())
111}
112
113#[cfg(test)]
114mod tests {
115    use arrow::array::Int32Array;
116    use arrow::datatypes::{DataType, Field, Schema};
117
118    use crate::batcher::table::batch::*;
119
120    fn batch(name: &str, values: Vec<i32>) -> RecordBatch {
121        RecordBatch::try_new(
122            Arc::new(Schema::new(vec![Field::new(name, DataType::Int32, false)])),
123            vec![Arc::new(Int32Array::from(values))],
124        )
125        .unwrap()
126    }
127
128    #[test]
129    fn test_combine_preserves_rows_and_rejects_schema_changes() {
130        let first = batch("a", vec![1, 2]);
131        let second = batch("a", vec![3]);
132        assert_eq!(
133            batch("a", vec![1, 2, 3]),
134            combine_batches([&first, &second]).unwrap()
135        );
136        let different = batch("b", vec![4]);
137        assert!(combine_batches([&first, &different]).is_err());
138        let schema = first.schema();
139        let metadata = Schema::new_with_metadata(
140            schema.fields().clone(),
141            [("version".to_string(), "2".to_string())]
142                .into_iter()
143                .collect(),
144        );
145        let changed = RecordBatch::try_new(Arc::new(metadata), first.columns().to_vec()).unwrap();
146        assert!(combine_batches([&first, &changed]).is_err());
147    }
148
149    #[test]
150    fn test_affected_rows_must_match_before_attribution() {
151        for (expected, actual, valid) in [(0, 0, true), (3, 3, true), (3, 2, false), (3, 4, false)]
152        {
153            assert_eq!(valid, validate_affected_rows(expected, actual).is_ok());
154        }
155    }
156}