risingwavelabs/risingwave · error · SinkError::Coordinator

should get metadata on checkpoint barrier

Error message

should get metadata on checkpoint barrier

What it means

On a checkpoint barrier the coordinator calls sink_writer.barrier(true), which must return metadata describing the sink's state. A None result means the writer did not produce metadata for a forced checkpoint, violating the coordinator protocol, so it wraps this as a SinkError::Coordinator.

Source

Thrown at src/connector/src/sink/coordinate.rs:215

                    schema_change,
                } => {
                    let prev_epoch = match state {
                        LogConsumerState::EpochBegun { curr_epoch } => curr_epoch,
                        _ => unreachable!("epoch must have begun before handling barrier"),
                    };
                    if is_checkpoint {
                        current_checkpoint += 1;
                        if current_checkpoint >= commit_checkpoint_interval.get()
                            || should_force_commit_on_checkpoint_barrier(
                                new_vnode_bitmap.is_some(),
                                is_stop,
                                schema_change.is_some(),
                            )
                        {
                            let start_time = Instant::now();
                            let metadata = sink_writer.barrier(true).await?;
                            let metadata = metadata.ok_or_else(|| {
                                SinkError::Coordinator(anyhow!(
                                    "should get metadata on checkpoint barrier"
                                ))
                            })?;
                            if schema_change.is_some() {
                                tracing::info!(
                                    sink_id = %self.param.sink_id,
                                    ?schema_change,
                                    "schema change received for coordinated log sinker"
                                );
                                assert!(
                                    is_stop,
                                    "schema change should stop current sink for sink {}",
                                    self.param.sink_id
                                );
                            }
                            coordinator_stream_handle
                                .commit(epoch, metadata, schema_change)
                                .await?;

View on GitHub (pinned to 6469eb736d)

Solutions

  1. Inspect the sink writer's barrier(true) implementation to ensure it always returns Some(metadata) on forced checkpoints.
  2. Verify the sink connector type supports coordination metadata (not all sinks do).
  3. Upgrade RisingWave / connector crate if this is a known sink implementation bug.
Defensive patterns

Strategy: try-catch

Try / catch

let metadata = sink_writer.barrier(true).await?
    .ok_or_else(|| SinkError::Coordinator(anyhow!("should get metadata on checkpoint barrier")))?;

Prevention

When it happens

Trigger: consume_log_and_sink processes a checkpoint barrier (with or without schema change) and `sink_writer.barrier(true).await` resolves to None.

Common situations: Using a sink writer implementation that does not support checkpoint metadata; bugs in a custom/dynamic sink writer where the checkpoint path returns None instead of Some(metadata).

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/c9b35de96c6ba1ce. Report an issue: GitHub.