1use std::collections::HashSet;
16use std::sync::Arc;
17use std::time::Instant;
18
19use catalog::CatalogManagerRef;
20use common_batcher::flush_limiter::FlushLimiter;
21use common_meta::node_manager::NodeManagerRef;
22use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
23use common_telemetry::{debug, warn};
24use partition::manager::PartitionRuleManagerRef;
25use session::context::QueryContextRef;
26use snafu::OptionExt;
27
28use crate::batcher::flow_notifier::FlowNotifier;
29use crate::batcher::logical_table::PHYSICAL_TABLE_KEY;
30use crate::batcher::logical_table::batch_convert::{
31 TableBatch, concat_modified_batches, transform_logical_batches_to_physical,
32};
33use crate::batcher::logical_table::flow_notifier::enqueue_flow_notifications;
34use crate::batcher::logical_table::pending_worker::FlushWaiter;
35use crate::batcher::logical_table::region_write::{
36 CatalogManagerPhysicalFlushAdapter, NodeManagerPhysicalFlushAdapter,
37 PartitionManagerPhysicalFlushAdapter, PhysicalFlushCatalogProvider, PhysicalFlushNodeRequester,
38 PhysicalFlushPartitionProvider, encode_region_write_requests, flush_region_writes_concurrently,
39 plan_region_batches, resolve_region_targets,
40};
41use crate::error;
42use crate::error::Result;
43use crate::metrics::{
44 FLUSH_DROPPED_ROWS, FLUSH_ELAPSED, FLUSH_FAILURES, FLUSH_ROWS, FLUSH_TOTAL,
45 PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED,
46};
47
48pub(in crate::batcher::logical_table) struct Batch {
49 pub(in crate::batcher::logical_table) table_batches: Vec<TableBatch>,
50 pub(in crate::batcher::logical_table) total_row_count: usize,
51 pub(in crate::batcher::logical_table) db_string: String,
52 pub(in crate::batcher::logical_table) ctx: QueryContextRef,
53 pub(in crate::batcher::logical_table) waiters: Vec<FlushWaiter>,
54}
55
56pub(in crate::batcher::logical_table) async fn spawn_flush(
57 flush: Batch,
58 partition_manager: PartitionRuleManagerRef,
59 node_manager: NodeManagerRef,
60 catalog_manager: CatalogManagerRef,
61 flow_notification_tx: FlowNotifier,
62 flush_limiter: FlushLimiter,
63) {
64 match flush_limiter.acquire().await {
65 Ok(permit) => {
66 tokio::spawn(async move {
67 let _permit = permit;
68 flush_batch_with_managers(
69 flush,
70 partition_manager,
71 node_manager,
72 catalog_manager,
73 flow_notification_tx,
74 )
75 .await;
76 });
77 }
78 Err(err) => {
79 warn!(err; "Flush semaphore closed, flushing inline");
80 flush_batch_with_managers(
81 flush,
82 partition_manager,
83 node_manager,
84 catalog_manager,
85 flow_notification_tx,
86 )
87 .await;
88 }
89 }
90}
91
92pub(in crate::batcher::logical_table) async fn flush_batch_with_managers(
93 flush: Batch,
94 partition_manager: PartitionRuleManagerRef,
95 node_manager: NodeManagerRef,
96 catalog_manager: CatalogManagerRef,
97 flow_notification_tx: FlowNotifier,
98) {
99 let partition_provider = PartitionManagerPhysicalFlushAdapter { partition_manager };
100 let node_requester = NodeManagerPhysicalFlushAdapter {
101 node_manager: node_manager.clone(),
102 };
103 let catalog_provider = CatalogManagerPhysicalFlushAdapter { catalog_manager };
104 flush_batch(
105 flush,
106 &partition_provider,
107 &node_requester,
108 &catalog_provider,
109 flow_notification_tx,
110 )
111 .await;
112}
113
114pub(in crate::batcher::logical_table) async fn flush_batch(
115 flush: Batch,
116 partition_manager: &(impl PhysicalFlushPartitionProvider + ?Sized),
117 node_manager: &(impl PhysicalFlushNodeRequester + ?Sized),
118 catalog_manager: &(impl PhysicalFlushCatalogProvider + ?Sized),
119 flow_notification_tx: FlowNotifier,
120) {
121 let Batch {
122 table_batches,
123 total_row_count,
124 db_string,
125 ctx,
126 waiters,
127 } = flush;
128 let start = Instant::now();
129
130 let physical_table_name = ctx
133 .extension(PHYSICAL_TABLE_KEY)
134 .unwrap_or(GREPTIME_PHYSICAL_TABLE)
135 .to_string();
136 let result = flush_batch_physical(
137 &table_batches,
138 &physical_table_name,
139 &ctx,
140 partition_manager,
141 node_manager,
142 catalog_manager,
143 )
144 .await;
145
146 let elapsed = start.elapsed().as_secs_f64();
147 FLUSH_ELAPSED.observe(elapsed);
148
149 debug!(
150 "Pending rows batch flushed, total rows: {}, elapsed time: {}s",
151 total_row_count, elapsed
152 );
153
154 match result {
155 Ok(affected_rows) => {
156 FLUSH_TOTAL.inc();
157 FLUSH_ROWS.observe(total_row_count as f64);
158 operator::metrics::DIST_INGEST_ROW_COUNT
159 .with_label_values(&[db_string.as_str()])
160 .inc_by(affected_rows as u64);
161
162 notify_waiters(waiters, Ok(()));
163 enqueue_flow_notifications(table_batches, &flow_notification_tx);
164 }
165 Err(err) => {
166 FLUSH_FAILURES.inc();
167 FLUSH_DROPPED_ROWS.inc_by(total_row_count as u64);
168 notify_waiters(waiters, Err(err));
169 }
170 }
171}
172
173pub async fn flush_batch_physical(
184 table_batches: &[TableBatch],
185 physical_table_name: &str,
186 ctx: &QueryContextRef,
187 partition_manager: &(impl PhysicalFlushPartitionProvider + ?Sized),
188 node_manager: &(impl PhysicalFlushNodeRequester + ?Sized),
189 catalog_manager: &(impl PhysicalFlushCatalogProvider + ?Sized),
190) -> Result<usize> {
191 let physical_table = {
193 let _timer = PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED
194 .with_label_values(&["flush_physical_resolve_table"])
195 .start_timer();
196 catalog_manager
197 .physical_table(
198 ctx.current_catalog(),
199 &ctx.current_schema(),
200 physical_table_name,
201 ctx.as_ref(),
202 )
203 .await?
204 .with_context(|| error::InternalSnafu {
205 err_msg: format!(
206 "Physical table '{}' not found during pending flush",
207 physical_table_name
208 ),
209 })?
210 };
211
212 let physical_table_info = physical_table.table_info;
213 let name_to_ids = physical_table
214 .col_name_to_ids
215 .with_context(|| error::InternalSnafu {
216 err_msg: format!(
217 "Physical table '{}' has no column IDs for pending flush",
218 physical_table_name
219 ),
220 })?;
221
222 let partition_rule = {
224 let _timer = PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED
225 .with_label_values(&["flush_physical_fetch_partition_rule"])
226 .start_timer();
227 partition_manager
228 .find_table_partition_rule(physical_table_info.as_ref())
229 .await?
230 };
231 let partition_columns = partition_rule.partition_columns();
232 let partition_columns_set: HashSet<&str> =
233 partition_columns.iter().map(String::as_str).collect();
234
235 let modified_batches =
237 transform_logical_batches_to_physical(table_batches, &name_to_ids, &partition_columns_set)?;
238
239 let combined_batch = concat_modified_batches(&modified_batches)?;
241
242 let physical_table_id = physical_table_info.table_id();
244 let planned_batches = plan_region_batches(
245 combined_batch,
246 physical_table_id,
247 partition_rule.as_ref(),
248 partition_columns,
249 )?;
250
251 let resolved_batches = resolve_region_targets(planned_batches, partition_manager).await?;
252 let region_writes = encode_region_write_requests(resolved_batches, ctx.skip_wal())?;
253 flush_region_writes_concurrently(node_manager, region_writes).await
254}
255
256pub(in crate::batcher::logical_table) fn notify_waiters(
257 waiters: Vec<FlushWaiter>,
258 result: Result<()>,
259) {
260 let shared_result = result.map_err(Arc::new);
261 for waiter in waiters {
262 let _ = waiter.response_tx.send(match &shared_result {
263 Ok(()) => Ok(()),
264 Err(error) => Err(Arc::clone(error)),
265 });
266 }
268}
269
270#[cfg(test)]
271mod tests {
272 use std::collections::HashMap;
273 use std::future::poll_fn;
274 use std::sync::Arc;
275 use std::sync::atomic::{AtomicUsize, Ordering};
276 use std::task::Poll;
277 use std::time::Duration;
278
279 use api::region::RegionResponse;
280 use api::v1::meta::Peer;
281 use api::v1::region::RegionRequest;
282 use arrow::record_batch::RecordBatch;
283 use async_trait::async_trait;
284 use catalog::error::Result as CatalogResult;
285 use common_batcher::request_limiter::RequestLimiter;
286 use common_meta::cache::TableFlownodeSetCacheRef;
287 use common_meta::node_manager::NodeManagerRef;
288 use datatypes::schema::{ColumnSchema as DtColumnSchema, Schema as DtSchema};
289 use partition::error::Result as PartitionResult;
290 use partition::partition::{PartitionRule, PartitionRuleRef, RegionMask};
291 use store_api::storage::RegionId;
292 use table::metadata::TableId;
293 use table::test_util::table_info::test_table_info;
294 use tokio::sync::{mpsc, oneshot};
295
296 use crate::batcher::flow_notifier::{FlowNotifier, start_flow_notification_worker};
297 use crate::batcher::logical_table::batch::{
298 Batch, flush_batch, flush_batch_physical, notify_waiters,
299 };
300 use crate::batcher::logical_table::batch_convert::TableBatch;
301 use crate::batcher::logical_table::pending_worker::FlushWaiter;
302 use crate::batcher::logical_table::region_write::{
303 PhysicalFlushCatalogProvider, PhysicalFlushNodeRequester, PhysicalFlushPartitionProvider,
304 PhysicalTableMetadata,
305 };
306 use crate::batcher::logical_table::test_util::{
307 FlowNotificationMockNodeManager, RecordingFlownode, mock_aligned_tag_batch,
308 };
309 use crate::batcher::test_util::mock_table_flownode_cache;
310 use crate::error;
311 use crate::error::Error;
312 use crate::metrics::FLOW_NOTIFICATION_DROPPED;
313
314 #[tokio::test]
315 async fn test_flush_batch_notifies_flownode_after_successful_physical_write() {
316 let table_id = 42;
317 let peer = Peer {
318 id: 7,
319 addr: "flow-7".to_string(),
320 };
321 let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
322 let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
323 let _requests_tx = requests_tx.clone();
324 let flow_node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
325 flownode: Arc::new(RecordingFlownode { requests_tx }),
326 });
327 let flow_notification_tx = mock_flow_notification_sender(cache, flow_node_manager.clone());
328 let ctx = session::context::QueryContext::arc();
329 let table_batches = vec![TableBatch {
330 table_name: "cpu".to_string(),
331 table_id,
332 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
333 row_count: 1,
334 }];
335 let writes = Arc::new(AtomicUsize::new(0));
336
337 flush_batch(
338 Batch {
339 table_batches,
340 total_row_count: 1,
341 db_string: ctx.get_db_string(),
342 ctx,
343 waiters: Vec::new(),
344 },
345 &MockFlushPartitionProvider {
346 partition_rule_calls: Arc::new(AtomicUsize::new(0)),
347 region_leader_calls: Arc::new(AtomicUsize::new(0)),
348 },
349 &MockFlushNodeRequester {
350 writes: writes.clone(),
351 fail: false,
352 },
353 &MockFlushCatalogProvider {
354 table: Some(mock_physical_table_metadata(1024)),
355 },
356 flow_notification_tx,
357 )
358 .await;
359
360 assert_eq!(1, writes.load(Ordering::SeqCst));
361 let requests = tokio::time::timeout(Duration::from_secs(1), requests_rx.recv())
362 .await
363 .unwrap()
364 .unwrap();
365 assert_eq!(
366 vec![api::v1::flow::DirtyWindowRequest {
367 table_id,
368 timestamps: vec![1000],
369 time_ranges: Vec::new(),
370 }],
371 requests.requests
372 );
373 }
374
375 #[tokio::test]
376 async fn test_flush_batch_does_not_notify_flownode_after_physical_write_error() {
377 let table_id = 42;
378 let peer = Peer {
379 id: 7,
380 addr: "flow-7".to_string(),
381 };
382 let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
383 let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
384 let _requests_tx = requests_tx.clone();
385 let flow_node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
386 flownode: Arc::new(RecordingFlownode { requests_tx }),
387 });
388 let flow_notification_tx = mock_flow_notification_sender(cache, flow_node_manager.clone());
389 let ctx = session::context::QueryContext::arc();
390 let table_batches = vec![TableBatch {
391 table_name: "cpu".to_string(),
392 table_id,
393 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
394 row_count: 1,
395 }];
396 let writes = Arc::new(AtomicUsize::new(0));
397
398 flush_batch(
399 Batch {
400 table_batches,
401 total_row_count: 1,
402 db_string: ctx.get_db_string(),
403 ctx,
404 waiters: Vec::new(),
405 },
406 &MockFlushPartitionProvider {
407 partition_rule_calls: Arc::new(AtomicUsize::new(0)),
408 region_leader_calls: Arc::new(AtomicUsize::new(0)),
409 },
410 &MockFlushNodeRequester {
411 writes: writes.clone(),
412 fail: true,
413 },
414 &MockFlushCatalogProvider {
415 table: Some(mock_physical_table_metadata(1024)),
416 },
417 flow_notification_tx,
418 )
419 .await;
420
421 assert_eq!(1, writes.load(Ordering::SeqCst));
422 assert!(
423 tokio::time::timeout(Duration::from_millis(50), requests_rx.recv())
424 .await
425 .is_err()
426 );
427 }
428
429 #[tokio::test]
430 async fn test_cancelled_waiter_retains_request_slot_until_notification() {
431 let limiter = RequestLimiter::try_new(1).unwrap();
432 let (response_tx, response_rx) = oneshot::channel();
433 let waiter = FlushWaiter {
434 response_tx,
435 _permit: limiter.acquire().await.unwrap(),
436 };
437 drop(response_rx);
438 let next = limiter.acquire();
439 tokio::pin!(next);
440 assert!(poll_fn(|cx| Poll::Ready(next.as_mut().poll(cx).is_pending())).await);
441 notify_waiters(vec![waiter], Ok(()));
442 let _permit = next.await.unwrap();
443 }
444
445 #[tokio::test]
446 async fn test_flush_batch_physical_uses_mockable_trait_dependencies() {
447 let table_batches = vec![TableBatch {
448 table_name: "t1".to_string(),
449 table_id: 11,
450 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
451 row_count: 1,
452 }];
453 let partition_calls = Arc::new(AtomicUsize::new(0));
454 let leader_calls = Arc::new(AtomicUsize::new(0));
455 let node = MockFlushNodeRequester::default();
456 let ctx = session::context::QueryContext::arc();
457
458 flush_batch_physical(
459 &table_batches,
460 "phy",
461 &ctx,
462 &MockFlushPartitionProvider {
463 partition_rule_calls: partition_calls.clone(),
464 region_leader_calls: leader_calls.clone(),
465 },
466 &node,
467 &MockFlushCatalogProvider {
468 table: Some(mock_physical_table_metadata(1024)),
469 },
470 )
471 .await
472 .unwrap();
473
474 assert_eq!(1, partition_calls.load(Ordering::SeqCst));
475 assert_eq!(1, leader_calls.load(Ordering::SeqCst));
476 assert_eq!(1, node.writes.load(Ordering::SeqCst));
477 }
478
479 #[tokio::test]
480 async fn test_flush_batch_physical_returns_actual_affected_rows() {
481 let table_batches = vec![TableBatch {
482 table_name: "t1".to_string(),
483 table_id: 11,
484 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
485 row_count: 1,
486 }];
487 let ctx = session::context::QueryContext::arc();
488
489 let affected_rows = flush_batch_physical(
490 &table_batches,
491 "phy",
492 &ctx,
493 &MockFlushPartitionProvider {
494 partition_rule_calls: Arc::new(AtomicUsize::new(0)),
495 region_leader_calls: Arc::new(AtomicUsize::new(0)),
496 },
497 &AffectedRowsFlushNodeRequester { affected_rows: 7 },
498 &MockFlushCatalogProvider {
499 table: Some(mock_physical_table_metadata(1024)),
500 },
501 )
502 .await
503 .unwrap();
504
505 assert_eq!(7, affected_rows);
506 }
507
508 #[tokio::test]
509 async fn test_flush_batch_physical_stops_before_partition_and_node_when_table_missing() {
510 let table_batches = vec![TableBatch {
511 table_name: "t1".to_string(),
512 table_id: 11,
513 batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
514 row_count: 1,
515 }];
516 let partition_calls = Arc::new(AtomicUsize::new(0));
517 let leader_calls = Arc::new(AtomicUsize::new(0));
518 let node = MockFlushNodeRequester::default();
519 let ctx = session::context::QueryContext::arc();
520
521 let err = flush_batch_physical(
522 &table_batches,
523 "missing_phy",
524 &ctx,
525 &MockFlushPartitionProvider {
526 partition_rule_calls: partition_calls.clone(),
527 region_leader_calls: leader_calls.clone(),
528 },
529 &node,
530 &MockFlushCatalogProvider { table: None },
531 )
532 .await
533 .unwrap_err();
534
535 assert!(
536 err.to_string()
537 .contains("Physical table 'missing_phy' not found")
538 );
539 assert_eq!(0, partition_calls.load(Ordering::SeqCst));
540 assert_eq!(0, leader_calls.load(Ordering::SeqCst));
541 assert_eq!(0, node.writes.load(Ordering::SeqCst));
542 }
543
544 #[tokio::test]
545 async fn test_flush_batch_physical_aborts_immediately_on_transform_error() {
546 let table_batches = vec![
547 TableBatch {
548 table_name: "broken".to_string(),
549 table_id: 11,
550 batches: vec![mock_aligned_tag_batch("unknown_tag", "host-1", 1000, 1.0)],
551 row_count: 1,
552 },
553 TableBatch {
554 table_name: "healthy".to_string(),
555 table_id: 12,
556 batches: vec![mock_aligned_tag_batch("tag1", "host-2", 2000, 2.0)],
557 row_count: 1,
558 },
559 ];
560 let partition_calls = Arc::new(AtomicUsize::new(0));
561 let leader_calls = Arc::new(AtomicUsize::new(0));
562 let node = MockFlushNodeRequester::default();
563 let ctx = session::context::QueryContext::arc();
564
565 let err = flush_batch_physical(
566 &table_batches,
567 "phy",
568 &ctx,
569 &MockFlushPartitionProvider {
570 partition_rule_calls: partition_calls.clone(),
571 region_leader_calls: leader_calls.clone(),
572 },
573 &node,
574 &MockFlushCatalogProvider {
575 table: Some(mock_physical_table_metadata(1024)),
576 },
577 )
578 .await
579 .unwrap_err();
580
581 assert!(err.to_string().contains("unknown_tag"));
582 assert_eq!(1, partition_calls.load(Ordering::SeqCst));
583 assert_eq!(0, leader_calls.load(Ordering::SeqCst));
584 assert_eq!(0, node.writes.load(Ordering::SeqCst));
585 }
586
587 fn mock_physical_table_metadata(table_id: TableId) -> PhysicalTableMetadata {
588 let schema = Arc::new(
589 DtSchema::try_new(vec![
590 DtColumnSchema::new(
591 "__primary_key",
592 datatypes::prelude::ConcreteDataType::binary_datatype(),
593 false,
594 ),
595 DtColumnSchema::new(
596 "greptime_timestamp",
597 datatypes::prelude::ConcreteDataType::timestamp_millisecond_datatype(),
598 false,
599 ),
600 DtColumnSchema::new(
601 "greptime_value",
602 datatypes::prelude::ConcreteDataType::float64_datatype(),
603 true,
604 ),
605 DtColumnSchema::new(
606 "tag1",
607 datatypes::prelude::ConcreteDataType::string_datatype(),
608 true,
609 ),
610 ])
611 .unwrap(),
612 );
613 let mut table_info = test_table_info(table_id, "phy", "public", "greptime", schema);
614 table_info.meta.column_ids = vec![0, 1, 2, 3];
615
616 PhysicalTableMetadata {
617 table_info: Arc::new(table_info),
618 col_name_to_ids: Some(HashMap::from([("tag1".to_string(), 3)])),
619 }
620 }
621
622 struct MockFlushCatalogProvider {
623 table: Option<PhysicalTableMetadata>,
624 }
625
626 #[async_trait]
627 impl PhysicalFlushCatalogProvider for MockFlushCatalogProvider {
628 async fn physical_table(
629 &self,
630 _catalog: &str,
631 _schema: &str,
632 _table_name: &str,
633 _query_ctx: &session::context::QueryContext,
634 ) -> CatalogResult<Option<PhysicalTableMetadata>> {
635 Ok(self.table.clone())
636 }
637 }
638
639 struct SingleRegionPartitionRule;
640
641 impl PartitionRule for SingleRegionPartitionRule {
642 fn as_any(&self) -> &dyn std::any::Any {
643 self
644 }
645
646 fn partition_columns(&self) -> &[String] {
647 &[]
648 }
649
650 fn find_region(
651 &self,
652 _values: &[datatypes::prelude::Value],
653 ) -> partition::error::Result<store_api::storage::RegionNumber> {
654 unimplemented!()
655 }
656
657 fn split_record_batch(
658 &self,
659 record_batch: &RecordBatch,
660 ) -> partition::error::Result<HashMap<store_api::storage::RegionNumber, RegionMask>>
661 {
662 Ok(HashMap::from([(
663 1,
664 RegionMask::new(
665 arrow::array::BooleanArray::from(vec![true; record_batch.num_rows()]),
666 record_batch.num_rows(),
667 ),
668 )]))
669 }
670 }
671
672 struct MockFlushPartitionProvider {
673 partition_rule_calls: Arc<AtomicUsize>,
674 region_leader_calls: Arc<AtomicUsize>,
675 }
676
677 #[async_trait]
678 impl PhysicalFlushPartitionProvider for MockFlushPartitionProvider {
679 async fn find_table_partition_rule(
680 &self,
681 _table_info: &table::metadata::TableInfo,
682 ) -> PartitionResult<PartitionRuleRef> {
683 self.partition_rule_calls.fetch_add(1, Ordering::SeqCst);
684 Ok(Arc::new(SingleRegionPartitionRule))
685 }
686
687 async fn find_region_leader(&self, _region_id: RegionId) -> error::Result<Peer> {
688 self.region_leader_calls.fetch_add(1, Ordering::SeqCst);
689 Ok(Peer {
690 id: 1,
691 addr: "node-1".to_string(),
692 })
693 }
694 }
695
696 #[derive(Default)]
697 struct MockFlushNodeRequester {
698 writes: Arc<AtomicUsize>,
699 fail: bool,
700 }
701
702 #[async_trait]
703 impl PhysicalFlushNodeRequester for MockFlushNodeRequester {
704 async fn handle(
705 &self,
706 _peer: &Peer,
707 _request: RegionRequest,
708 ) -> error::Result<RegionResponse> {
709 self.writes.fetch_add(1, Ordering::SeqCst);
710 if self.fail {
711 return Err(Error::Internal {
712 err_msg: "physical write failed".to_string(),
713 });
714 }
715 Ok(RegionResponse::new(0))
716 }
717 }
718
719 fn mock_flow_notification_sender(
720 cache: TableFlownodeSetCacheRef,
721 node_manager: NodeManagerRef,
722 ) -> FlowNotifier {
723 let (tx, rx) = FlowNotifier::try_new(16, FLOW_NOTIFICATION_DROPPED.clone()).unwrap();
724 start_flow_notification_worker(rx, cache, node_manager);
725 tx
726 }
727
728 #[derive(Default)]
729 struct AffectedRowsFlushNodeRequester {
730 affected_rows: usize,
731 }
732
733 #[async_trait]
734 impl PhysicalFlushNodeRequester for AffectedRowsFlushNodeRequester {
735 async fn handle(
736 &self,
737 _peer: &Peer,
738 _request: RegionRequest,
739 ) -> error::Result<RegionResponse> {
740 Ok(RegionResponse::new(self.affected_rows))
741 }
742 }
743}