risingwavelabs/risingwave · error · SinkError::Remote

epoch has not been initialize, call `begin_epoch`

Error message

epoch has not been initialize, call `begin_epoch`

What it means

The remote sink writer requires an epoch to be set via begin_epoch before any data can be written. write_batch finds self.epoch is None, meaning a chunk arrived without a preceding BeginEpoch call in the writer's lifecycle, violating the per-epoch write protocol.

Solutions

  1. Ensure begin_epoch(epoch) is called on the writer before the first write_batch/write_string for each epoch.
  2. In tests, call writer.begin_epoch(epoch).await before writing chunks.
  3. Check the recovery path replays the epoch-begin request before data chunks.

Example fix

// before
writer.write_batch(chunk).await?;
// after
writer.begin_epoch(epoch).await?;
writer.write_batch(chunk).await?;
Defensive patterns

Strategy: type-guard

Type guard

fn ensure_epoch_ready(writer: &RemoteSink) -> Result<u64, SinkError> {
    writer.epoch.ok_or_else(|| SinkError::Remote(anyhow!("epoch has not been initialize, call `begin_epoch`")))
}

Try / catch

match writer.write_batch(chunk).await {
    Err(e) if e.to_string().contains("epoch has not been initialize") => {
        writer.begin_epoch(epoch).await?;
        writer.write_batch(chunk).await?;
    }
    other => other,
}

Prevention

When it happens

Trigger: Calling write_batch (e.g. from test_remote_sink or the stream executor) on a RemoteSink writer where begin_epoch(epoch) was never invoked, or after a state where epoch was reset/never initialized.

Common situations: Test harness forgetting to call begin_epoch before write_batch; recovery path instantiating the writer and writing chunks before replaying the epoch-begin; misuse of the SinkWriter API by an embedder.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/remote.rs:633

            batch_id: 0,
            stream_handle,
            metrics: SinkWriterMetrics::for_test(),
        }
    }
}

#[async_trait]
impl SinkWriter for CoordinatedRemoteSinkWriter {
    type CommitMetadata = Option<SinkMetadata>;

    async fn write_batch(&mut self, chunk: StreamChunk) -> Result<()> {
        let cardinality = chunk.cardinality();
        self.metrics
            .connector_sink_rows_received
            .inc_by(cardinality as _);

        let epoch = self.epoch.ok_or_else(|| {
            SinkError::Remote(anyhow!("epoch has not been initialize, call `begin_epoch`"))
        })?;
        let batch_id = self.batch_id;
        self.stream_handle
            .request_sender
            .send_request(JniSinkWriterStreamRequest::Chunk {
                chunk,
                epoch,
                batch_id,
            })
            .await?;
        self.batch_id += 1;
        Ok(())
    }

    async fn begin_epoch(&mut self, epoch: u64) -> Result<()> {
        self.epoch = Some(epoch);
        Ok(())
    }

View on GitHub (pinned to 6469eb736d)