1mod 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
80pub(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
94pub 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 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 pub(crate) async fn accepts_bulk_destinations(
209 &self,
210 batches: impl Iterator<Item = &(QueryContextRef, RowInsertRequests)>,
211 ) -> bool {
212 let mut checked = HashSet::new();
217 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 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 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 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}