1use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
16use std::fmt::Debug;
17use std::future::{Future, poll_fn};
18use std::pin::pin;
19use std::sync::LazyLock;
20use std::task::Poll;
21use std::time::{Duration, Instant};
22
23use anyhow::anyhow;
24use await_tree::InstrumentAwait;
25use futures::future::{Either, pending, select};
26use futures::pin_mut;
27use itertools::Itertools;
28use risingwave_common::bail;
29use risingwave_common::bitmap::Bitmap;
30use risingwave_common::util::retry::exponential_backoff;
31use risingwave_connector::connector_common::IcebergSinkCompactionUpdate;
32use risingwave_connector::dispatch_sink;
33use risingwave_connector::sink::catalog::SinkId;
34use risingwave_connector::sink::{
35 Sink, SinkCommitCoordinator, SinkCommittedEpochSubscriber, SinkError, SinkParam, build_sink,
36};
37use risingwave_meta_model::pending_sink_state::SinkState;
38use risingwave_pb::connector_service::{SinkMetadata, coordinate_request};
39use risingwave_pb::stream_plan::PbSinkSchemaChange;
40use sea_orm::DatabaseConnection;
41use thiserror_ext::AsReport;
42use tokio::select;
43use tokio::sync::Semaphore;
44use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
45use tokio::time::sleep;
46use tokio_retry::strategy::jitter;
47use tonic::Status;
48use tracing::{error, warn};
49
50use crate::manager::exactly_once_util::{
51 clean_aborted_records, commit_and_prune_epoch, list_sink_states_ordered_by_epoch,
52 persist_pre_commit_metadata,
53};
54use crate::manager::sink_coordination::handle::SinkWriterCoordinationHandle;
55
56const MAX_CONCURRENT_COORDINATOR_INITIALIZATIONS: usize = 8;
59static COORDINATOR_INIT_SEMAPHORE: LazyLock<Semaphore> =
60 LazyLock::new(|| Semaphore::new(MAX_CONCURRENT_COORDINATOR_INITIALIZATIONS));
61
62async fn run_future_with_periodic_fn<F: Future>(
63 future: F,
64 interval: Duration,
65 mut f: impl FnMut(),
66) -> F::Output {
67 pin_mut!(future);
68 loop {
69 match select(&mut future, pin!(sleep(interval))).await {
70 Either::Left((output, _)) => {
71 break output;
72 }
73 Either::Right(_) => f(),
74 }
75 }
76}
77
78type HandleId = usize;
79
80#[derive(Default)]
81struct AligningRequests<R> {
82 requests: Vec<R>,
83 handle_ids: HashSet<HandleId>,
84 committed_bitmap: Option<Bitmap>, }
86
87impl<R> AligningRequests<R> {
88 fn add_new_request(
89 &mut self,
90 handle_id: HandleId,
91 request: R,
92 vnode_bitmap: &Bitmap,
93 ) -> anyhow::Result<()>
94 where
95 R: Debug,
96 {
97 let committed_bitmap = self
98 .committed_bitmap
99 .get_or_insert_with(|| Bitmap::zeros(vnode_bitmap.len()));
100 assert_eq!(committed_bitmap.len(), vnode_bitmap.len());
101
102 let check_bitmap = (&*committed_bitmap) & vnode_bitmap;
103 if check_bitmap.count_ones() > 0 {
104 return Err(anyhow!(
105 "duplicate vnode {:?}. request vnode: {:?}, prev vnode: {:?}. pending request: {:?}, request: {:?}",
106 check_bitmap.iter_ones().collect_vec(),
107 vnode_bitmap,
108 committed_bitmap,
109 self.requests,
110 request
111 ));
112 }
113 *committed_bitmap |= vnode_bitmap;
114 self.requests.push(request);
115 assert!(self.handle_ids.insert(handle_id));
116 Ok(())
117 }
118
119 fn aligned(&self) -> bool {
120 self.committed_bitmap.as_ref().is_some_and(|b| b.all())
121 }
122}
123
124type RetryBackoffFuture = std::pin::Pin<Box<tokio::time::Sleep>>;
125type RetryBackoffStrategy = impl Iterator<Item = RetryBackoffFuture> + Send + 'static;
126
127struct TwoPhaseCommitHandler {
128 db: DatabaseConnection,
129 sink_id: SinkId,
130 curr_hummock_committed_epoch: u64,
131 job_committed_epoch_rx: UnboundedReceiver<u64>,
132 last_committed_epoch: Option<u64>,
133 pending_epochs: VecDeque<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)>,
134 prepared_epochs: VecDeque<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)>,
135 backoff_state: Option<(RetryBackoffFuture, RetryBackoffStrategy)>,
136}
137
138impl TwoPhaseCommitHandler {
139 fn new(
140 db: DatabaseConnection,
141 sink_id: SinkId,
142 initial_hummock_committed_epoch: u64,
143 job_committed_epoch_rx: UnboundedReceiver<u64>,
144 last_committed_epoch: Option<u64>,
145 ) -> Self {
146 Self {
147 db,
148 sink_id,
149 curr_hummock_committed_epoch: initial_hummock_committed_epoch,
150 job_committed_epoch_rx,
151 last_committed_epoch,
152 pending_epochs: VecDeque::new(),
153 prepared_epochs: VecDeque::new(),
154 backoff_state: None,
155 }
156 }
157
158 #[define_opaque(RetryBackoffStrategy)]
159 fn get_retry_backoff_strategy() -> RetryBackoffStrategy {
160 exponential_backoff(Duration::from_millis(10), 10, Duration::from_secs(60))
161 .map(jitter)
162 .map(|delay| Box::pin(tokio::time::sleep(delay)))
163 }
164
165 async fn next_to_commit(
166 &mut self,
167 ) -> anyhow::Result<(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>)> {
168 loop {
169 let wait_backoff = async {
170 if self.prepared_epochs.is_empty() {
171 pending::<()>().await;
172 } else if let Some((backoff_fut, _)) = &mut self.backoff_state {
173 backoff_fut.await;
174 }
175 };
176
177 select! {
178 _ = wait_backoff => {
179 let item = self.prepared_epochs.front().cloned().expect("non-empty");
180 return Ok(item);
181 }
182
183 recv_epoch = self.job_committed_epoch_rx.recv() => {
184 let Some(recv_epoch) = recv_epoch else {
185 return Err(anyhow!(
186 "Hummock committed epoch sender closed unexpectedly"
187 ));
188 };
189 self.curr_hummock_committed_epoch = recv_epoch;
190 while let Some((epoch, metadata, schema_change)) = self.pending_epochs.pop_front_if(|(epoch, _, _)| *epoch <= recv_epoch) {
191 if let Some((last_epoch, _, _)) = self.prepared_epochs.back() {
192 assert!(epoch > *last_epoch, "prepared epochs must be in increasing order");
193 }
194 self.prepared_epochs.push_back((epoch, metadata, schema_change));
195 }
196 }
197 }
198 }
199 }
200
201 fn push_new_item(
202 &mut self,
203 epoch: u64,
204 metadata: Option<Vec<u8>>,
205 schema_change: Option<PbSinkSchemaChange>,
206 ) {
207 if epoch > self.curr_hummock_committed_epoch {
208 if let Some((last_epoch, _, _)) = self.pending_epochs.back() {
209 assert!(
210 epoch > *last_epoch,
211 "pending epochs must be in increasing order"
212 );
213 }
214 self.pending_epochs
215 .push_back((epoch, metadata, schema_change));
216 } else {
217 assert!(self.pending_epochs.is_empty());
218 if let Some((last_epoch, _, _)) = self.prepared_epochs.back() {
219 assert!(
220 epoch > *last_epoch,
221 "prepared epochs must be in increasing order"
222 );
223 }
224 self.prepared_epochs
225 .push_back((epoch, metadata, schema_change));
226 }
227 }
228
229 async fn ack_committed(&mut self, epoch: u64) -> anyhow::Result<()> {
230 self.backoff_state = None;
231 let (last_epoch, _, _) = self.prepared_epochs.pop_front().expect("non-empty");
232 assert_eq!(last_epoch, epoch);
233
234 commit_and_prune_epoch(&self.db, self.sink_id, epoch, self.last_committed_epoch).await?;
235 self.last_committed_epoch = Some(epoch);
236 Ok(())
237 }
238
239 fn failed_committed(&mut self, epoch: u64, err: SinkError) {
240 assert_eq!(self.prepared_epochs.front().expect("non-empty").0, epoch,);
241 if let Some((prev_fut, strategy)) = &mut self.backoff_state {
242 let new_fut = strategy.next().expect("infinite");
243 *prev_fut = new_fut;
244 } else {
245 let mut strategy = Self::get_retry_backoff_strategy();
246 let backoff_fut = strategy.next().expect("infinite");
247 self.backoff_state = Some((backoff_fut, strategy));
248 }
249 tracing::error!(
250 error = %err.as_report(),
251 %self.sink_id,
252 "failed to commit epoch {}, Retrying after backoff",
253 epoch,
254 );
255 }
256
257 fn is_empty(&self) -> bool {
258 self.pending_epochs.is_empty() && self.prepared_epochs.is_empty()
259 }
260
261 fn has_uncommitted_schema_change(&self) -> bool {
266 if let Some((_, _, schema_change)) = self.pending_epochs.back() {
267 schema_change.is_some()
268 } else if let Some((_, _, schema_change)) = self.prepared_epochs.back() {
269 schema_change.is_some()
270 } else {
271 false
272 }
273 }
274}
275
276struct CoordinationHandleManager {
277 param: SinkParam,
278 writer_handles: HashMap<HandleId, SinkWriterCoordinationHandle>,
279 next_handle_id: HandleId,
280 request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
281}
282
283impl CoordinationHandleManager {
284 fn start(
285 &mut self,
286 log_store_rewind_start_epoch: Option<u64>,
287 handle_ids: impl IntoIterator<Item = HandleId>,
288 ) -> anyhow::Result<()> {
289 for handle_id in handle_ids {
290 let handle = self
291 .writer_handles
292 .get_mut(&handle_id)
293 .ok_or_else(|| anyhow!("failed to find handle {} to start", handle_id,))?;
294 handle.start(log_store_rewind_start_epoch).map_err(|_| {
295 anyhow!(
296 "failed to start {:?} for handle {}",
297 log_store_rewind_start_epoch,
298 handle_id
299 )
300 })?;
301 }
302 Ok(())
303 }
304
305 fn ack_aligned_initial_epoch(&mut self, aligned_initial_epoch: u64) -> anyhow::Result<()> {
306 for (handle_id, handle) in &mut self.writer_handles {
307 handle
308 .ack_aligned_initial_epoch(aligned_initial_epoch)
309 .map_err(|_| {
310 anyhow!(
311 "failed to ack aligned initial epoch {:?} for handle {}",
312 aligned_initial_epoch,
313 handle_id
314 )
315 })?;
316 }
317 Ok(())
318 }
319
320 fn ack_commit(
321 &mut self,
322 epoch: u64,
323 handle_ids: impl IntoIterator<Item = HandleId>,
324 ) -> anyhow::Result<()> {
325 for handle_id in handle_ids {
326 let handle = self.writer_handles.get_mut(&handle_id).ok_or_else(|| {
327 anyhow!(
328 "failed to find handle {} when acknowledging the commit for epoch {}",
329 handle_id,
330 epoch
331 )
332 })?;
333 handle.ack_commit(epoch).map_err(|_| {
334 anyhow!(
335 "failed to acknowledge the commit for epoch {} on handle {}",
336 epoch,
337 handle_id
338 )
339 })?;
340 }
341 Ok(())
342 }
343
344 async fn next_request_inner(
345 writer_handles: &mut HashMap<HandleId, SinkWriterCoordinationHandle>,
346 ) -> anyhow::Result<(HandleId, coordinate_request::Msg)> {
347 poll_fn(|cx| {
348 for (handle_id, handle) in writer_handles.iter_mut() {
349 if let Poll::Ready(result) = handle.poll_next_request(cx) {
350 return Poll::Ready(result.map(|request| (*handle_id, request)));
351 }
352 }
353 Poll::Pending
354 })
355 .await
356 }
357}
358
359enum CoordinationHandleManagerEvent {
360 NewHandle,
361 UpdateVnodeBitmap,
362 Stop,
363 CommitRequest {
364 epoch: u64,
365 metadata: SinkMetadata,
366 schema_change: Option<PbSinkSchemaChange>,
367 },
368 AlignInitialEpoch(u64),
369}
370
371impl CoordinationHandleManagerEvent {
372 fn name(&self) -> &'static str {
373 match self {
374 CoordinationHandleManagerEvent::NewHandle => "NewHandle",
375 CoordinationHandleManagerEvent::UpdateVnodeBitmap => "UpdateVnodeBitmap",
376 CoordinationHandleManagerEvent::Stop => "Stop",
377 CoordinationHandleManagerEvent::CommitRequest { .. } => "CommitRequest",
378 CoordinationHandleManagerEvent::AlignInitialEpoch(_) => "AlignInitialEpoch",
379 }
380 }
381}
382
383impl CoordinationHandleManager {
384 async fn next_event(&mut self) -> anyhow::Result<(HandleId, CoordinationHandleManagerEvent)> {
385 select! {
386 handle = self.request_rx.recv() => {
387 let handle = handle.ok_or_else(|| anyhow!("end of writer request stream"))?;
388 if handle.param() != &self.param {
389 warn!(prev_param = ?self.param, new_param = ?handle.param(), "sink param mismatch");
390 }
391 let handle_id = self.next_handle_id;
392 self.next_handle_id += 1;
393 self.writer_handles.insert(handle_id, handle);
394 Ok((handle_id, CoordinationHandleManagerEvent::NewHandle))
395 }
396 result = Self::next_request_inner(&mut self.writer_handles) => {
397 let (handle_id, request) = result?;
398 let event = match request {
399 coordinate_request::Msg::CommitRequest(request) => {
400 CoordinationHandleManagerEvent::CommitRequest {
401 epoch: request.epoch,
402 metadata: request.metadata.ok_or_else(|| anyhow!("empty sink metadata"))?,
403 schema_change: request.schema_change,
404 }
405 }
406 coordinate_request::Msg::AlignInitialEpochRequest(epoch) => {
407 CoordinationHandleManagerEvent::AlignInitialEpoch(epoch)
408 }
409 coordinate_request::Msg::UpdateVnodeRequest(_) => {
410 CoordinationHandleManagerEvent::UpdateVnodeBitmap
411 }
412 coordinate_request::Msg::Stop(_) => {
413 CoordinationHandleManagerEvent::Stop
414 }
415 coordinate_request::Msg::StartRequest(_) => {
416 unreachable!("should have been handled");
417 }
418 };
419 Ok((handle_id, event))
420 }
421 }
422 }
423
424 fn vnode_bitmap(&self, handle_id: HandleId) -> &Bitmap {
425 self.writer_handles[&handle_id].vnode_bitmap()
426 }
427
428 fn stop_handle(&mut self, handle_id: HandleId) -> anyhow::Result<()> {
429 self.writer_handles
430 .remove(&handle_id)
431 .expect("should exist")
432 .stop()
433 }
434
435 async fn wait_init_handles(&mut self) -> anyhow::Result<HashSet<HandleId>> {
436 assert!(self.writer_handles.is_empty());
437 let mut init_requests = AligningRequests::default();
438 while !init_requests.aligned() {
439 let (handle_id, event) = self.next_event().await?;
440 let unexpected_event = match event {
441 CoordinationHandleManagerEvent::NewHandle => {
442 init_requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
443 continue;
444 }
445 event => event.name(),
446 };
447 return Err(anyhow!(
448 "expect new handle during init, but got {}",
449 unexpected_event
450 ));
451 }
452 Ok(init_requests.handle_ids)
453 }
454
455 async fn alter_parallelisms(
456 &mut self,
457 altered_handles: impl Iterator<Item = HandleId>,
458 ) -> anyhow::Result<HashSet<HandleId>> {
459 let mut requests = AligningRequests::default();
460 for handle_id in altered_handles {
461 requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
462 }
463 let mut remaining_handles: HashSet<_> = self
464 .writer_handles
465 .keys()
466 .filter(|handle_id| !requests.handle_ids.contains(handle_id))
467 .cloned()
468 .collect();
469 while !remaining_handles.is_empty() || !requests.aligned() {
470 let (handle_id, event) = self.next_event().await?;
471 match event {
472 CoordinationHandleManagerEvent::NewHandle => {
473 requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
474 }
475 CoordinationHandleManagerEvent::UpdateVnodeBitmap => {
476 assert!(remaining_handles.remove(&handle_id));
477 requests.add_new_request(handle_id, (), self.vnode_bitmap(handle_id))?;
478 }
479 CoordinationHandleManagerEvent::Stop => {
480 assert!(remaining_handles.remove(&handle_id));
481 self.stop_handle(handle_id)?;
482 }
483 CoordinationHandleManagerEvent::CommitRequest { epoch, .. } => {
484 bail!(
485 "receive commit request on epoch {} from handle {} during alter parallelism",
486 epoch,
487 handle_id
488 );
489 }
490 CoordinationHandleManagerEvent::AlignInitialEpoch(epoch) => {
491 bail!(
492 "receive AlignInitialEpoch on epoch {} from handle {} during alter parallelism",
493 epoch,
494 handle_id
495 );
496 }
497 }
498 }
499 Ok(requests.handle_ids)
500 }
501}
502
503enum CoordinatorWorkerState {
509 Running,
510 WaitingForFlushed(HashSet<HandleId>),
511}
512
513pub struct CoordinatorWorker {
514 handle_manager: CoordinationHandleManager,
515 last_writer_acked_epoch: Option<u64>,
521 curr_state: CoordinatorWorkerState,
522}
523
524enum CoordinatorWorkerEvent {
525 HandleManagerEvent(HandleId, CoordinationHandleManagerEvent),
526 ReadyToCommit(u64, Option<Vec<u8>>, Option<PbSinkSchemaChange>),
527}
528
529impl CoordinatorWorker {
530 pub async fn run(
531 param: SinkParam,
532 request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
533 db: DatabaseConnection,
534 subscriber: SinkCommittedEpochSubscriber,
535 iceberg_compact_stat_sender: UnboundedSender<IcebergSinkCompactionUpdate>,
536 ) {
537 let sink = match build_sink(param.clone()) {
538 Ok(sink) => sink,
539 Err(e) => {
540 error!(
541 error = %e.as_report(),
542 "unable to build sink with param {:?}",
543 param
544 );
545 return;
546 }
547 };
548
549 dispatch_sink!(sink, sink, {
550 let coordinator =
551 match Box::pin(sink.new_coordinator(Some(iceberg_compact_stat_sender))).await {
552 Ok(coordinator) => coordinator,
553 Err(e) => {
554 error!(
555 error = %e.as_report(),
556 "unable to build coordinator with param {:?}",
557 param
558 );
559 return;
560 }
561 };
562 Self::execute_coordinator(db, param, request_rx, coordinator, subscriber).await
563 });
564 }
565
566 pub async fn execute_coordinator(
567 db: DatabaseConnection,
568 param: SinkParam,
569 request_rx: UnboundedReceiver<SinkWriterCoordinationHandle>,
570 coordinator: SinkCommitCoordinator,
571 subscriber: SinkCommittedEpochSubscriber,
572 ) {
573 let mut worker = CoordinatorWorker {
574 handle_manager: CoordinationHandleManager {
575 param,
576 writer_handles: HashMap::new(),
577 next_handle_id: 0,
578 request_rx,
579 },
580 last_writer_acked_epoch: None,
581 curr_state: CoordinatorWorkerState::Running,
582 };
583
584 if let Err(e) = worker.run_coordination(db, coordinator, subscriber).await {
585 for handle in worker.handle_manager.writer_handles.into_values() {
586 handle.abort(Status::internal(format!(
587 "failed to run coordination: {:?}",
588 e.as_report()
589 )))
590 }
591 }
592 }
593
594 async fn try_handle_init_requests(
595 &mut self,
596 pending_handle_ids: &HashSet<HandleId>,
597 two_phase_handler: &mut TwoPhaseCommitHandler,
598 ) -> anyhow::Result<()> {
599 assert!(matches!(self.curr_state, CoordinatorWorkerState::Running));
600 if two_phase_handler.has_uncommitted_schema_change() {
601 self.curr_state = CoordinatorWorkerState::WaitingForFlushed(pending_handle_ids.clone());
603 } else {
604 self.handle_init_requests_impl(pending_handle_ids.clone())
605 .await?;
606 }
607 Ok(())
608 }
609
610 async fn handle_init_requests_impl(
611 &mut self,
612 pending_handle_ids: impl IntoIterator<Item = HandleId>,
613 ) -> anyhow::Result<()> {
614 let log_store_rewind_start_epoch = self.last_writer_acked_epoch;
615 self.handle_manager
616 .start(log_store_rewind_start_epoch, pending_handle_ids)?;
617 if log_store_rewind_start_epoch.is_none() {
618 let mut align_requests = AligningRequests::default();
619 while !align_requests.aligned() {
620 let (handle_id, event) = self.handle_manager.next_event().await?;
621 match event {
622 CoordinationHandleManagerEvent::AlignInitialEpoch(initial_epoch) => {
623 align_requests.add_new_request(
624 handle_id,
625 initial_epoch,
626 self.handle_manager.vnode_bitmap(handle_id),
627 )?;
628 }
629 other => {
630 return Err(anyhow!("expect AlignInitialEpoch but got {}", other.name()));
631 }
632 }
633 }
634 let aligned_initial_epoch = align_requests
635 .requests
636 .into_iter()
637 .max()
638 .expect("non-empty");
639 self.handle_manager
640 .ack_aligned_initial_epoch(aligned_initial_epoch)?;
641 }
642 Ok(())
643 }
644
645 async fn next_event(
646 &mut self,
647 two_phase_handler: &mut TwoPhaseCommitHandler,
648 ) -> anyhow::Result<CoordinatorWorkerEvent> {
649 if let CoordinatorWorkerState::WaitingForFlushed(pending_handle_ids) = &self.curr_state
650 && two_phase_handler.is_empty()
651 {
652 let pending_handle_ids = pending_handle_ids.clone();
653 self.handle_init_requests_impl(pending_handle_ids).await?;
654 self.curr_state = CoordinatorWorkerState::Running;
655 }
656
657 select! {
658 next_handle_event = self.handle_manager.next_event() => {
659 let (handle_id, event) = next_handle_event?;
660 Ok(CoordinatorWorkerEvent::HandleManagerEvent(handle_id, event))
661 }
662
663 next_item_to_commit = two_phase_handler.next_to_commit() => {
664 let (epoch, metadata, schema_change) = next_item_to_commit?;
665 Ok(CoordinatorWorkerEvent::ReadyToCommit(epoch, metadata, schema_change))
666 }
667 }
668 }
669
670 async fn run_coordination(
671 &mut self,
672 db: DatabaseConnection,
673 mut coordinator: SinkCommitCoordinator,
674 subscriber: SinkCommittedEpochSubscriber,
675 ) -> anyhow::Result<()> {
676 let sink_id = self.handle_manager.param.sink_id;
677
678 let coordinator_init_permit = COORDINATOR_INIT_SEMAPHORE
679 .acquire()
680 .instrument_await("acquire_sink_coordinator_init_permit")
681 .await?;
682 let mut two_phase_handler = self
683 .init_state_from_store(&db, sink_id, subscriber, &mut coordinator)
684 .await?;
685 match &mut coordinator {
686 SinkCommitCoordinator::SinglePhase(coordinator) => coordinator.init().await?,
687 SinkCommitCoordinator::TwoPhase(coordinator) => coordinator.init().await?,
688 }
689 drop(coordinator_init_permit);
690
691 let mut running_handles = self.handle_manager.wait_init_handles().await?;
692 self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
693 .await?;
694
695 let mut pending_epochs: BTreeMap<u64, AligningRequests<_>> = BTreeMap::new();
696 let mut pending_new_handles = vec![];
697 loop {
698 let event = self.next_event(&mut two_phase_handler).await?;
699 let (handle_id, epoch, commit_request) = match event {
700 CoordinatorWorkerEvent::HandleManagerEvent(handle_id, event) => match event {
701 CoordinationHandleManagerEvent::NewHandle => {
702 pending_new_handles.push(handle_id);
703 continue;
704 }
705 CoordinationHandleManagerEvent::UpdateVnodeBitmap => {
706 running_handles = self
707 .handle_manager
708 .alter_parallelisms(pending_new_handles.drain(..).chain([handle_id]))
709 .await?;
710 self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
711 .await?;
712 continue;
713 }
714 CoordinationHandleManagerEvent::Stop => {
715 self.handle_manager.stop_handle(handle_id)?;
716 running_handles = self
717 .handle_manager
718 .alter_parallelisms(pending_new_handles.drain(..))
719 .await?;
720 self.try_handle_init_requests(&running_handles, &mut two_phase_handler)
721 .await?;
722
723 continue;
724 }
725 CoordinationHandleManagerEvent::CommitRequest {
726 epoch,
727 metadata,
728 schema_change,
729 } => (handle_id, epoch, (metadata, schema_change)),
730 CoordinationHandleManagerEvent::AlignInitialEpoch(_) => {
731 bail!("receive AlignInitialEpoch after initialization")
732 }
733 },
734 CoordinatorWorkerEvent::ReadyToCommit(epoch, metadata, schema_change) => {
735 let start_time = Instant::now();
736 let commit_fut = async {
737 match &mut coordinator {
738 SinkCommitCoordinator::SinglePhase(coordinator) => {
739 assert!(metadata.is_none());
740 if let Some(schema_change) = schema_change {
741 coordinator
742 .commit_schema_change(epoch, schema_change)
743 .instrument_await(Self::commit_span(
744 "single_phase_schema_change",
745 sink_id,
746 epoch,
747 ))
748 .await?;
749 }
750 }
751 SinkCommitCoordinator::TwoPhase(coordinator) => {
752 if let Some(metadata) = metadata {
753 coordinator
754 .commit_data(epoch, metadata)
755 .instrument_await(Self::commit_span(
756 "two_phase_commit_data",
757 sink_id,
758 epoch,
759 ))
760 .await?;
761 }
762 if let Some(schema_change) = schema_change {
763 coordinator
764 .commit_schema_change(epoch, schema_change)
765 .instrument_await(Self::commit_span(
766 "two_phase_commit_schema_change",
767 sink_id,
768 epoch,
769 ))
770 .await?;
771 }
772 }
773 }
774 Ok(())
775 };
776 let commit_res =
777 run_future_with_periodic_fn(commit_fut, Duration::from_secs(5), || {
778 warn!(
779 elapsed = ?start_time.elapsed(),
780 %sink_id,
781 "committing"
782 );
783 })
784 .await;
785
786 match commit_res {
787 Ok(_) => {
788 two_phase_handler.ack_committed(epoch).await?;
789 }
790 Err(e) => {
791 two_phase_handler.failed_committed(epoch, e);
792 }
793 }
794
795 continue;
796 }
797 };
798 if !running_handles.contains(&handle_id) {
799 bail!(
800 "receiving commit request from non-running handle {}, running handles: {:?}",
801 handle_id,
802 running_handles
803 );
804 }
805 pending_epochs.entry(epoch).or_default().add_new_request(
806 handle_id,
807 commit_request,
808 self.handle_manager.vnode_bitmap(handle_id),
809 )?;
810 if pending_epochs
811 .first_key_value()
812 .expect("non-empty")
813 .1
814 .aligned()
815 {
816 let (epoch, commit_requests) = pending_epochs.pop_first().expect("non-empty");
817 let mut metadatas = Vec::with_capacity(commit_requests.requests.len());
818 let mut requests = commit_requests.requests.into_iter();
819 let (first_metadata, first_schema_change) = requests.next().expect("non-empty");
820 metadatas.push(first_metadata);
821 for (metadata, schema_change) in requests {
822 if first_schema_change != schema_change {
823 return Err(anyhow!(
824 "got different schema change {:?} to prev schema change {:?}",
825 schema_change,
826 first_schema_change
827 ));
828 }
829 metadatas.push(metadata);
830 }
831
832 match &mut coordinator {
833 SinkCommitCoordinator::SinglePhase(coordinator) => {
834 if !metadatas.is_empty() {
835 let start_time = Instant::now();
836 run_future_with_periodic_fn(
837 coordinator.commit_data(epoch, metadatas).instrument_await(
838 Self::commit_span("single_phase_commit_data", sink_id, epoch),
839 ),
840 Duration::from_secs(5),
841 || {
842 warn!(
843 elapsed = ?start_time.elapsed(),
844 %sink_id,
845 "committing"
846 );
847 },
848 )
849 .await
850 .map_err(|e| anyhow!(e))?;
851 }
852 if first_schema_change.is_some() {
853 persist_pre_commit_metadata(
854 &db,
855 sink_id as _,
856 epoch,
857 None,
858 first_schema_change.as_ref(),
859 )
860 .await?;
861 two_phase_handler.push_new_item(epoch, None, first_schema_change);
862 }
863 }
864 SinkCommitCoordinator::TwoPhase(coordinator) => {
865 let commit_metadata = coordinator
866 .pre_commit(epoch, metadatas, first_schema_change.clone())
867 .instrument_await(Self::commit_span(
868 "two_phase_pre_commit",
869 sink_id,
870 epoch,
871 ))
872 .await?;
873 persist_pre_commit_metadata(
879 &db,
880 sink_id as _,
881 epoch,
882 commit_metadata.clone(),
883 first_schema_change.as_ref(),
884 )
885 .await?;
886 two_phase_handler.push_new_item(
887 epoch,
888 commit_metadata,
889 first_schema_change,
890 );
891 }
892 }
893
894 self.handle_manager
895 .ack_commit(epoch, commit_requests.handle_ids)?;
896 self.last_writer_acked_epoch = Some(epoch);
897 }
898 }
899 }
900
901 async fn init_state_from_store(
903 &mut self,
904 db: &DatabaseConnection,
905 sink_id: SinkId,
906 subscriber: SinkCommittedEpochSubscriber,
907 coordinator: &mut SinkCommitCoordinator,
908 ) -> anyhow::Result<TwoPhaseCommitHandler> {
909 let ordered_metadata = list_sink_states_ordered_by_epoch(db, sink_id as _).await?;
910
911 let mut metadata_iter = ordered_metadata.into_iter().peekable();
912 let last_committed_epoch = metadata_iter
913 .next_if(|(_, state, _, _)| matches!(state, SinkState::Committed))
914 .map(|(epoch, _, _, _)| epoch);
915
916 let pending_items = metadata_iter
917 .peeking_take_while(|(_, state, _, _)| matches!(state, SinkState::Pending))
918 .map(|(epoch, _, metadata, schema_change)| (epoch, metadata, schema_change))
919 .collect_vec();
920 self.last_writer_acked_epoch = pending_items
921 .last()
922 .map(|(epoch, _, _)| *epoch)
923 .or(last_committed_epoch);
924
925 let mut aborted_epochs = vec![];
926
927 for (epoch, state, metadata, _) in metadata_iter {
928 match state {
929 SinkState::Aborted => {
930 if let Some(metadata) = metadata
931 && let SinkCommitCoordinator::TwoPhase(coordinator) = coordinator
932 {
933 coordinator.abort(epoch, metadata).await;
934 }
935 aborted_epochs.push(epoch);
936 }
937 other => {
938 unreachable!(
939 "unexpected state {:?} after pending items at epoch {}",
940 other, epoch
941 );
942 }
943 }
944 }
945
946 clean_aborted_records(db, sink_id, aborted_epochs).await?;
948
949 let (initial_hummock_committed_epoch, job_committed_epoch_rx) = subscriber(sink_id).await?;
950 let mut two_phase_handler = TwoPhaseCommitHandler::new(
951 db.clone(),
952 sink_id,
953 initial_hummock_committed_epoch,
954 job_committed_epoch_rx,
955 last_committed_epoch,
956 );
957
958 for (epoch, metadata, schema_change) in pending_items {
959 two_phase_handler.push_new_item(epoch, metadata, schema_change);
960 }
961
962 Ok(two_phase_handler)
963 }
964
965 fn commit_span(stage: &str, sink_id: SinkId, epoch: u64) -> await_tree::Span {
966 await_tree::span!("sink_coord_{stage} (sink_id {sink_id}, epoch {epoch})").long_running()
967 }
968}