1use 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 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 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
61pub(in crate::batcher::logical_table) struct TableResolutionPlan {
64 pub(in crate::batcher::logical_table) region_schemas: HashMap<String, (Arc<ArrowSchema>, u32)>,
66 pub(in crate::batcher::logical_table) tables_to_create: Vec<(String, Vec<ColumnSchema>)>,
68 pub(in crate::batcher::logical_table) tables_to_alter: Vec<(String, Vec<String>)>,
70}
71
72impl LogicalTablePendingRowsBatcher {
73 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 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
124fn 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 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 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 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 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 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 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 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 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}