{"record":{"id":"cb8bcf3d5da1a7a5","repo":"risingwavelabs/risingwave","slug":"epoch-has-not-been-initialize-call-begin-epoch","errorCode":null,"errorMessage":"epoch has not been initialize, call `begin_epoch`","messagePattern":"epoch has not been initialize, call `begin_epoch`","errorType":"exception","errorClass":"SinkError::Remote","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":633,"sourceCode":"            batch_id: 0,\n            stream_handle,\n            metrics: SinkWriterMetrics::for_test(),\n        }\n    }\n}\n\n#[async_trait]\nimpl SinkWriter for CoordinatedRemoteSinkWriter {\n    type CommitMetadata = Option<SinkMetadata>;\n\n    async fn write_batch(&mut self, chunk: StreamChunk) -> Result<()> {\n        let cardinality = chunk.cardinality();\n        self.metrics\n            .connector_sink_rows_received\n            .inc_by(cardinality as _);\n\n        let epoch = self.epoch.ok_or_else(|| {\n            SinkError::Remote(anyhow!(\"epoch has not been initialize, call `begin_epoch`\"))\n        })?;\n        let batch_id = self.batch_id;\n        self.stream_handle\n            .request_sender\n            .send_request(JniSinkWriterStreamRequest::Chunk {\n                chunk,\n                epoch,\n                batch_id,\n            })\n            .await?;\n        self.batch_id += 1;\n        Ok(())\n    }\n\n    async fn begin_epoch(&mut self, epoch: u64) -> Result<()> {\n        self.epoch = Some(epoch);\n        Ok(())\n    }","sourceCodeStart":615,"sourceCodeEnd":651,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L615-L651","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// before\nwriter.write_batch(chunk).await?;\n// after\nwriter.begin_epoch(epoch).await?;\nwriter.write_batch(chunk).await?;","handlingStrategy":"type-guard","validationCode":null,"typeGuard":"fn ensure_epoch_ready(writer: &RemoteSink) -> Result<u64, SinkError> {\n    writer.epoch.ok_or_else(|| SinkError::Remote(anyhow!(\"epoch has not been initialize, call `begin_epoch`\")))\n}","tryCatchPattern":"match writer.write_batch(chunk).await {\n    Err(e) if e.to_string().contains(\"epoch has not been initialize\") => {\n        writer.begin_epoch(epoch).await?;\n        writer.write_batch(chunk).await?;\n    }\n    other => other,\n}","preventionTips":["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"],"tags":["sink","epoch","lifecycle","api-misuse"],"backgroundTag":"invalid-state-transition","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}