servers/batcher/table/
batch.rs1use 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
33pub(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 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}