Skip to main content

servers/batcher/logical_table/
tables.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::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use api::v1::{ColumnSchema, RowInsertRequests, Rows, SemanticType};
19use arrow::datatypes::Schema as ArrowSchema;
20use async_trait::async_trait;
21use common_time::timestamp::TimeUnit;
22use session::context::QueryContextRef;
23use snafu::OptionExt;
24
25use crate::batcher::logical_table::LogicalTablePendingRowsBatcher;
26use crate::batcher::logical_table::batch_convert::RecordBatchWithTsIdx;
27use crate::error;
28use crate::error::Result;
29use crate::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED;
30use crate::prom_row_builder::{
31    build_prom_create_table_schema_from_proto, identify_missing_columns_from_proto,
32    rows_to_aligned_record_batch,
33};
34
35#[async_trait]
36pub trait PendingRowsSchemaAlterer: Send + Sync {
37    /// Batch-create multiple logical tables that are missing.
38    /// Each entry is `(table_name, request_schema)`.
39    async fn create_tables_if_missing_batch(
40        &self,
41        catalog: &str,
42        schema: &str,
43        tables: &[(&str, &[ColumnSchema])],
44        with_metric_engine: bool,
45        ctx: QueryContextRef,
46    ) -> Result<()>;
47
48    /// Batch-alter multiple logical tables to add missing tag columns.
49    /// Each entry is `(table_name, missing_column_names)`.
50    async fn add_missing_prom_tag_columns_batch(
51        &self,
52        catalog: &str,
53        schema: &str,
54        tables: &[(&str, &[String])],
55        ctx: QueryContextRef,
56    ) -> Result<()>;
57}
58
59pub type PendingRowsSchemaAltererRef = Arc<dyn PendingRowsSchemaAlterer>;
60
61/// Intermediate planning state for resolving and preparing logical tables
62/// before row-to-batch alignment.
63pub(in crate::batcher::logical_table) struct TableResolutionPlan {
64    /// Resolved table schema and table id by logical table name.
65    pub(in crate::batcher::logical_table) region_schemas: HashMap<String, (Arc<ArrowSchema>, u32)>,
66    /// Missing tables that need to be created before alignment.
67    pub(in crate::batcher::logical_table) tables_to_create: Vec<(String, Vec<ColumnSchema>)>,
68    /// Existing tables that need tag-column schema evolution.
69    pub(in crate::batcher::logical_table) tables_to_alter: Vec<(String, Vec<String>)>,
70}
71
72impl LogicalTablePendingRowsBatcher {
73    /// Converts proto `RowInsertRequests` directly into aligned `RecordBatch`es
74    /// in a single pass, handling table creation, schema alteration, column
75    /// renaming, reordering, and null-filling without building intermediate
76    /// RecordBatches.
77    pub(in crate::batcher::logical_table) async fn build_and_align_table_batches(
78        &self,
79        requests: &RowInsertRequests,
80        ctx: &QueryContextRef,
81    ) -> Result<(Vec<(String, u32, RecordBatchWithTsIdx)>, usize)> {
82        let catalog = ctx.current_catalog().to_string();
83        let schema = ctx.current_schema();
84
85        let (table_rows, total_rows) = Self::collect_non_empty_table_rows(requests);
86        if total_rows == 0 {
87            return Ok((Vec::new(), 0));
88        }
89
90        let unique_tables = Self::collect_unique_table_schemas(&table_rows)?;
91        let mut plan = self
92            .plan_table_resolution(&catalog, &schema, ctx, &unique_tables)
93            .await?;
94
95        // New tables are created on the request's selected physical table;
96        // their time index must use the physical table's unit (a missing
97        // physical table is auto-created as millisecond by the schema
98        // alterer, matching the default here).
99        if !plan.tables_to_create.is_empty() {
100            let physical_unit = self.physical_time_index_unit_or_default(ctx).await;
101            for (_, request_schema) in &mut plan.tables_to_create {
102                align_create_schema_time_index(request_schema, physical_unit);
103            }
104        }
105
106        self.create_missing_tables_and_refresh_schemas(
107            &catalog,
108            &schema,
109            ctx,
110            &table_rows,
111            &mut plan,
112        )
113        .await?;
114
115        self.alter_tables_and_refresh_schemas(&catalog, &schema, ctx, &mut plan)
116            .await?;
117
118        let aligned_batches = Self::build_aligned_batches(&table_rows, &plan.region_schemas)?;
119
120        Ok((aligned_batches, total_rows))
121    }
122}
123
124/// Rewrites the time index column of a create-table schema to `unit`, if the
125/// schema carries a timestamp column in another unit.
126fn align_create_schema_time_index(request_schema: &mut [ColumnSchema], unit: TimeUnit) {
127    for column in request_schema {
128        if column.semantic_type == SemanticType::Timestamp as i32 {
129            column.datatype = api::helper::timestamp_datatype(unit) as i32;
130            column.datatype_extension = None;
131        }
132    }
133}
134
135impl LogicalTablePendingRowsBatcher {
136    /// Extracts non-empty `(table_name, rows)` pairs and computes total row
137    /// count across the retained entries.
138    pub(in crate::batcher::logical_table) fn collect_non_empty_table_rows(
139        requests: &RowInsertRequests,
140    ) -> (Vec<(&str, &Rows)>, usize) {
141        let mut table_rows: Vec<(&str, &Rows)> = Vec::with_capacity(requests.inserts.len());
142        let mut total_rows = 0;
143
144        for request in &requests.inserts {
145            let Some(rows) = &request.rows else {
146                continue;
147            };
148            if rows.rows.is_empty() {
149                continue;
150            }
151
152            total_rows += rows.rows.len();
153            table_rows.push((request.table_name.as_str(), rows));
154        }
155
156        (table_rows, total_rows)
157    }
158}
159
160impl LogicalTablePendingRowsBatcher {
161    /// Returns unique `(table_name, proto_schema)` pairs while keeping the
162    /// first-seen schema for duplicate table names.
163    pub(in crate::batcher::logical_table) fn collect_unique_table_schemas<'a>(
164        table_rows: &[(&'a str, &'a Rows)],
165    ) -> Result<Vec<(&'a str, &'a [ColumnSchema])>> {
166        let mut unique_tables: Vec<(&str, &[ColumnSchema])> = Vec::with_capacity(table_rows.len());
167        let mut seen = HashSet::new();
168
169        for (table_name, rows) in table_rows {
170            if seen.insert(*table_name) {
171                unique_tables.push((*table_name, &rows.schema));
172            } else {
173                // table_rows should group rows by table name.
174                return error::InvalidPromRemoteRequestSnafu {
175                    msg: format!(
176                        "Found duplicated table name in RowInsertRequest: {}",
177                        table_name
178                    ),
179                }
180                .fail();
181            }
182        }
183
184        Ok(unique_tables)
185    }
186}
187
188impl LogicalTablePendingRowsBatcher {
189    /// Resolves table metadata and classifies each table into existing,
190    /// to-create, and to-alter groups used by subsequent DDL steps.
191    pub(in crate::batcher::logical_table) async fn plan_table_resolution(
192        &self,
193        catalog: &str,
194        schema: &str,
195        ctx: &QueryContextRef,
196        unique_tables: &[(&str, &[ColumnSchema])],
197    ) -> Result<TableResolutionPlan> {
198        let mut plan = TableResolutionPlan {
199            region_schemas: HashMap::with_capacity(unique_tables.len()),
200            tables_to_create: Vec::new(),
201            tables_to_alter: Vec::new(),
202        };
203
204        let resolved_tables = {
205            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
206                .with_label_values(&["align_resolve_table"])
207                .start_timer();
208            futures::future::join_all(unique_tables.iter().map(|(table_name, _)| {
209                self.catalog_manager
210                    .table(catalog, schema, table_name, Some(ctx.as_ref()))
211            }))
212            .await
213        };
214
215        for ((table_name, rows_schema), table_result) in unique_tables.iter().zip(resolved_tables) {
216            let table = table_result?;
217
218            if let Some(table) = table {
219                let table_info = table.table_info();
220                let table_id = table_info.ident.table_id;
221                let region_schema = table_info.meta.schema.arrow_schema().clone();
222
223                let missing_columns = {
224                    let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
225                        .with_label_values(&["align_identify_missing_columns"])
226                        .start_timer();
227                    identify_missing_columns_from_proto(rows_schema, region_schema.as_ref())?
228                };
229                if !missing_columns.is_empty() {
230                    plan.tables_to_alter
231                        .push(((*table_name).to_string(), missing_columns));
232                }
233                plan.region_schemas
234                    .insert((*table_name).to_string(), (region_schema, table_id));
235            } else {
236                let request_schema = {
237                    let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
238                        .with_label_values(&["align_build_create_table_schema"])
239                        .start_timer();
240                    build_prom_create_table_schema_from_proto(rows_schema)?
241                };
242                plan.tables_to_create
243                    .push(((*table_name).to_string(), request_schema));
244            }
245        }
246
247        Ok(plan)
248    }
249}
250
251impl LogicalTablePendingRowsBatcher {
252    /// Batch-creates missing tables, refreshes their schema metadata, and
253    /// enqueues follow-up alters for extra tag columns discovered in later rows.
254    pub(in crate::batcher::logical_table) async fn create_missing_tables_and_refresh_schemas(
255        &self,
256        catalog: &str,
257        schema: &str,
258        ctx: &QueryContextRef,
259        table_rows: &[(&str, &Rows)],
260        plan: &mut TableResolutionPlan,
261    ) -> Result<()> {
262        if plan.tables_to_create.is_empty() {
263            return Ok(());
264        }
265
266        let create_refs: Vec<(&str, &[ColumnSchema])> = plan
267            .tables_to_create
268            .iter()
269            .map(|(name, schema)| (name.as_str(), schema.as_slice()))
270            .collect();
271
272        {
273            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
274                .with_label_values(&["align_batch_create_tables"])
275                .start_timer();
276            self.schema_alterer
277                .create_tables_if_missing_batch(
278                    catalog,
279                    schema,
280                    &create_refs,
281                    self.prom_store_with_metric_engine,
282                    ctx.clone(),
283                )
284                .await?;
285        }
286
287        let created_table_results = {
288            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
289                .with_label_values(&["align_resolve_table_after_create"])
290                .start_timer();
291            futures::future::join_all(plan.tables_to_create.iter().map(|(table_name, _)| {
292                self.catalog_manager
293                    .table(catalog, schema, table_name, Some(ctx.as_ref()))
294            }))
295            .await
296        };
297
298        for ((table_name, _), table_result) in
299            plan.tables_to_create.iter().zip(created_table_results)
300        {
301            let table = table_result?.with_context(|| error::UnexpectedResultSnafu {
302                reason: format!(
303                    "Table not found after pending batch create attempt: {}",
304                    table_name
305                ),
306            })?;
307            let table_info = table.table_info();
308            let table_id = table_info.ident.table_id;
309            let region_schema = table_info.meta.schema.arrow_schema().clone();
310            plan.region_schemas
311                .insert(table_name.clone(), (region_schema, table_id));
312        }
313
314        Self::enqueue_alter_for_new_tables(table_rows, plan)?;
315
316        Ok(())
317    }
318}
319
320impl LogicalTablePendingRowsBatcher {
321    /// For newly created tables, re-checks all row schemas and appends alter
322    /// operations when additional tag columns are still missing.
323    pub(in crate::batcher::logical_table) fn enqueue_alter_for_new_tables(
324        table_rows: &[(&str, &Rows)],
325        plan: &mut TableResolutionPlan,
326    ) -> Result<()> {
327        let created_tables: HashSet<&str> = plan
328            .tables_to_create
329            .iter()
330            .map(|(table_name, _)| table_name.as_str())
331            .collect();
332
333        for (table_name, rows) in table_rows {
334            if !created_tables.contains(table_name) {
335                continue;
336            }
337
338            let Some((region_schema, _)) = plan.region_schemas.get(*table_name) else {
339                continue;
340            };
341
342            let missing_columns = identify_missing_columns_from_proto(&rows.schema, region_schema)?;
343            if missing_columns.is_empty()
344                || plan
345                    .tables_to_alter
346                    .iter()
347                    .any(|(existing_name, _)| existing_name == *table_name)
348            {
349                continue;
350            }
351
352            plan.tables_to_alter
353                .push((table_name.to_string(), missing_columns));
354        }
355
356        Ok(())
357    }
358}
359
360impl LogicalTablePendingRowsBatcher {
361    /// Batch-alters tables that have missing tag columns and refreshes the
362    /// in-memory schema map used for row alignment.
363    pub(in crate::batcher::logical_table) async fn alter_tables_and_refresh_schemas(
364        &self,
365        catalog: &str,
366        schema: &str,
367        ctx: &QueryContextRef,
368        plan: &mut TableResolutionPlan,
369    ) -> Result<()> {
370        if plan.tables_to_alter.is_empty() {
371            return Ok(());
372        }
373
374        let alter_refs: Vec<(&str, &[String])> = plan
375            .tables_to_alter
376            .iter()
377            .map(|(name, cols)| (name.as_str(), cols.as_slice()))
378            .collect();
379        {
380            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
381                .with_label_values(&["align_batch_add_missing_columns"])
382                .start_timer();
383            self.schema_alterer
384                .add_missing_prom_tag_columns_batch(catalog, schema, &alter_refs, ctx.clone())
385                .await?;
386        }
387
388        let altered_table_results = {
389            let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
390                .with_label_values(&["align_resolve_table_after_schema_alter"])
391                .start_timer();
392            futures::future::join_all(plan.tables_to_alter.iter().map(|(table_name, _)| {
393                self.catalog_manager
394                    .table(catalog, schema, table_name, Some(ctx.as_ref()))
395            }))
396            .await
397        };
398
399        for ((table_name, _), table_result) in
400            plan.tables_to_alter.iter().zip(altered_table_results)
401        {
402            let table = table_result?.with_context(|| error::UnexpectedResultSnafu {
403                reason: format!(
404                    "Table not found after pending batch schema alter: {}",
405                    table_name
406                ),
407            })?;
408            let table_info = table.table_info();
409            let table_id = table_info.ident.table_id;
410            let refreshed_region_schema = table_info.meta.schema.arrow_schema().clone();
411            plan.region_schemas
412                .insert(table_name.clone(), (refreshed_region_schema, table_id));
413        }
414
415        Ok(())
416    }
417}
418
419impl LogicalTablePendingRowsBatcher {
420    /// Converts proto rows to `RecordBatch` values aligned to resolved region
421    /// schemas and returns `(table_name, table_id, batch)` tuples.
422    pub(in crate::batcher::logical_table) fn build_aligned_batches(
423        table_rows: &[(&str, &Rows)],
424        region_schemas: &HashMap<String, (Arc<ArrowSchema>, u32)>,
425    ) -> Result<Vec<(String, u32, RecordBatchWithTsIdx)>> {
426        let mut aligned_batches = Vec::with_capacity(table_rows.len());
427        for (table_name, rows) in table_rows {
428            let (region_schema, table_id) =
429                region_schemas.get(*table_name).cloned().with_context(|| {
430                    error::UnexpectedResultSnafu {
431                        reason: format!("Region schema not resolved for table: {}", table_name),
432                    }
433                })?;
434
435            let record_batch = {
436                let _timer = PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED
437                    .with_label_values(&["align_rows_to_record_batch"])
438                    .start_timer();
439                rows_to_aligned_record_batch(rows, region_schema.as_ref())?
440            };
441            aligned_batches.push((table_name.to_string(), table_id, record_batch));
442        }
443
444        Ok(aligned_batches)
445    }
446}
447
448#[cfg(test)]
449mod tests {
450
451    use api::v1::value::ValueData;
452    use api::v1::{
453        ColumnDataType, ColumnSchema, Row, RowInsertRequest, RowInsertRequests, Rows, SemanticType,
454        Value,
455    };
456    use arrow::datatypes::{DataType as ArrowDataType, Field, Schema as ArrowSchema};
457    use common_query::prelude::greptime_timestamp;
458
459    use crate::batcher::logical_table::LogicalTablePendingRowsBatcher;
460    use crate::batcher::logical_table::batch_convert::TableBatch;
461    use crate::batcher::logical_table::flow_notifier::extract_timestamps;
462    use crate::prom_row_builder::rows_to_aligned_record_batch;
463
464    #[test]
465    fn test_extract_timestamps_uses_aligned_custom_timestamp_index() {
466        let rows = Rows {
467            schema: vec![
468                ColumnSchema {
469                    column_name: greptime_timestamp().to_string(),
470                    datatype: ColumnDataType::TimestampMillisecond as i32,
471                    semantic_type: SemanticType::Timestamp as i32,
472                    ..Default::default()
473                },
474                ColumnSchema {
475                    column_name: "host".to_string(),
476                    datatype: ColumnDataType::String as i32,
477                    semantic_type: SemanticType::Tag as i32,
478                    ..Default::default()
479                },
480                ColumnSchema {
481                    column_name: "greptime_value".to_string(),
482                    datatype: ColumnDataType::Float64 as i32,
483                    semantic_type: SemanticType::Field as i32,
484                    ..Default::default()
485                },
486            ],
487            rows: vec![
488                Row {
489                    values: vec![
490                        Value {
491                            value_data: Some(ValueData::TimestampMillisecondValue(1000)),
492                        },
493                        Value {
494                            value_data: Some(ValueData::StringValue("host-1".to_string())),
495                        },
496                        Value {
497                            value_data: Some(ValueData::F64Value(1.0)),
498                        },
499                    ],
500                },
501                Row {
502                    values: vec![
503                        Value {
504                            value_data: Some(ValueData::TimestampMillisecondValue(2000)),
505                        },
506                        Value {
507                            value_data: Some(ValueData::StringValue("host-2".to_string())),
508                        },
509                        Value {
510                            value_data: Some(ValueData::F64Value(2.0)),
511                        },
512                    ],
513                },
514            ],
515        };
516        let target_schema = ArrowSchema::new(vec![
517            Field::new("host", ArrowDataType::Utf8, true),
518            Field::new(
519                "timestamp",
520                ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
521                false,
522            ),
523            Field::new("greptime_value", ArrowDataType::Float64, true),
524        ]);
525        let batch = rows_to_aligned_record_batch(&rows, &target_schema).unwrap();
526        assert_eq!(1, batch.timestamp_index);
527        let table_batch = TableBatch {
528            table_name: "cpu".to_string(),
529            table_id: 42,
530            row_count: batch.batch.num_rows(),
531            batches: vec![batch],
532        };
533
534        assert_eq!(vec![1000, 2000], extract_timestamps(&table_batch));
535    }
536
537    #[test]
538    fn test_collect_non_empty_table_rows_filters_empty_payloads() {
539        let requests = RowInsertRequests {
540            inserts: vec![
541                RowInsertRequest {
542                    table_name: "cpu".to_string(),
543                    rows: Some(mock_rows(2, "host")),
544                },
545                RowInsertRequest {
546                    table_name: "mem".to_string(),
547                    rows: Some(mock_rows(0, "host")),
548                },
549                RowInsertRequest {
550                    table_name: "disk".to_string(),
551                    rows: None,
552                },
553            ],
554        };
555
556        let (table_rows, total_rows) =
557            LogicalTablePendingRowsBatcher::collect_non_empty_table_rows(&requests);
558
559        assert_eq!(2, total_rows);
560        assert_eq!(1, table_rows.len());
561        assert_eq!("cpu", table_rows[0].0);
562        assert_eq!(2, table_rows[0].1.rows.len());
563    }
564
565    fn mock_rows(row_count: usize, schema_name: &str) -> Rows {
566        Rows {
567            schema: vec![ColumnSchema {
568                column_name: schema_name.to_string(),
569                ..Default::default()
570            }],
571            rows: (0..row_count).map(|_| Row { values: vec![] }).collect(),
572        }
573    }
574}