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
- Ensure begin_epoch(epoch) is called on the writer before the first write_batch/write_string for each epoch.
- In tests, call writer.begin_epoch(epoch).await before writing chunks.
- 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
- Always call begin_epoch before the first write in every epoch, including after recovery
- Encode the epoch lifecycle in a small wrapper type that enforces begin-then-write
- Add unit tests that exercise begin_epoch -> write -> commit ordering
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
- get different response epoch to commit epoch
- newly start epoch after update vnode bitmap not matched…
- replace sink requires a sink job
- Actor exited unexpectedly
- all senders are dropped
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)