risingwavelabs/risingwave · error
historical read semaphore closed
Error message
historical read semaphore closed
What it means
KvLogStoreReader::start_from acquires a permit from a historical-read semaphore (limiting concurrent state-store reads). Acquire fails only when the semaphore is closed, meaning the semaphore's owner was dropped. The reader cannot proceed to read persisted log data from the state store.
Solutions
- Ensure the semaphore owner lives at least as long as every reader using it (hold an Arc in the reader).
- Check for an early drop of the reader/factory context during recovery or migration.
- Treat the error as terminal: recreate reader and semaphore pair instead of retrying acquire.
- In tests, keep the semaphore owner alive for the whole read duration.
Example fix
// before: semaphore owned by a dropped context let reader = KvLogStoreReader::new(ctx /* ctx owns semaphore */); drop(ctx); reader.start_from(...).await?; // after: reader holds its own Arc<Semaphore> let semaphore = Arc::new(Semaphore::new(n)); let reader = KvLogStoreReader::with_semaphore(semaphore.clone());
Defensive patterns
Strategy: try-catch
Validate before calling
if semaphore.is_closed() {
return Err("historical read semaphore already closed; recreate reader");
} Type guard
fn semaphore_live(s: &Arc<Semaphore>) -> bool { Arc::strong_count(s) > 1 } // owner still around Try / catch
match reader.start_from(...).await {
Err(e) if e.to_string().contains("historical read semaphore closed") => {
// recreate reader+semaphore pair
rebuild_reader_and_retry()
}
r => r,
} Prevention
- Store the semaphore in an Arc owned by the reader itself, not only by an outer context.
- Keep the owning context alive for the reader's full lifetime.
- Detect early teardown of the reader/factory during recovery.
When it happens
Trigger: Reader's start_from awaiting `semaphore.acquire_owned()` after the semaphore has been closed/dropped — its owning struct (e.g. the shared reader context or the actor that created it) is gone.
Common situations: Reader outliving its owning context (e.g. used after the streaming actor or factory was dropped); mis-managed Arc where the only semaphore owner was dropped during recovery; long-lived reader used in tests after teardown.
Related errors
- cannot get truncated epoch
- failed to receive update vnode
- should get the first epoch
- unable to send barrier
- unable to send init epoch
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/243fc060bbbc0eab.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/common/log_store_impl/kv_log_store/reader.rs:310
// This limits the number of concurrent historical reads across
// the compute node, preventing excessive state store I/O during
// recovery or restart.
//
// Deadlock safety: The writer side of this log store runs as a
// separate future in a `select` with this consumer. Even if this
// consumer is blocked waiting for a permit, the writer can still
// make progress processing barriers. The log store buffer is also
// unbounded on the writer side, so the writer won't be blocked by
// the reader not consuming. Therefore, blocking here does not
// prevent barrier flow.
let permit = if let Some(semaphore) = &self.historical_read_semaphore {
Some(
semaphore
.clone()
.acquire_owned()
.instrument_await("Wait for Historical Read Permit")
.await
.map_err(|_| anyhow!("historical read semaphore closed"))?,
)
} else {
None
};
KvLogStoreReaderFutureState::ReadStateStoreStream(
self.read_persisted_log_store(range_start).await?,
permit,
)
};
self.rx.rewind(start_offset);
Ok(())
}
async fn init(&mut self) -> LogStoreResult<()> {
if let Some(init_epoch_rx) = self.init_epoch_rx.take() {
let init_epoch = init_epoch_rx
.await
.map_err(|_| anyhow!("should get the first epoch"))?;View on GitHub (pinned to 6469eb736d)