risingwavelabs/risingwave · error
end of upstream input
Error message
end of upstream input
What it means
The sync KV log store's write path awaits a future that receives a write batch from upstream. If that upstream future completes with None (the upstream stream/input ended) instead of a write batch, next_event maps it to this error — the log store cannot continue writing without upstream input.
Solutions
- Find why the upstream write stream ended (component shutdown, earlier error, channel drop) in the logs.
- Ensure the log-store writer/upstream stays alive for the duration of writes; fix premature drop in code or tests.
- Retry the operation after the log store / actor is recovered.
Example fix
// before (test) drop(upstream_tx); write_future.await; // -> "end of upstream input" // after upstream_tx.send(write_batch).unwrap(); // keep upstream alive until writes finish
Defensive patterns
Strategy: fallback
Validate before calling
// before writing, confirm the upstream writer stream is still open
if upstream_tx.is_closed() { return Err(upstream_gone()); } Type guard
fn upstream_alive(tx: &UnboundedSender<WriteBatch>) -> bool { !tx.is_closed() } Try / catch
match next_event().await {
Err(e) if e.to_string().contains("end of upstream input") => {
// upstream closed mid-write: flush what we have and reopen the log store writer
reopen_writer_and_retry();
}
other => other,
} Prevention
- Keep the upstream write channel open until all pending writes complete
- Sequence shutdown: finish writes before dropping the log store writer
- In tests, send a WriteBatch instead of dropping the upstream sender
When it happens
Trigger: WriteFuture::ReceiveFromUpstream resolves with opt = None, i.e. the upstream stream delivering write batches terminated while a write operation was still pending.
Common situations: Upstream log/writer component shut down or errored while a write was in flight, Hummock/log-store teardown racing a pending write, or madsim tests dropping the upstream channel.
Related errors
- a stream has reached the end but some other stream has not…
- cannot get truncated epoch
- Exchange executor should not have children!
- Filter can only receive bool array
- should get the first epoch
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/be43be4a68985e15.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/sync_kv_log_store.rs:519
sleep_future,
..
} => {
if let Some(sleep_future) = sleep_future {
sleep_future.await;
metrics
.pause_duration_ns
.inc_by(start_instant.elapsed().as_nanos() as _);
tracing::trace!("resuming write future");
}
must_match!(replace(self, WriteFuture::Empty), WriteFuture::Paused { stream, write_state, barrier, .. } => {
Ok((stream, write_state, WriteFutureEvent::UpstreamMessageReceived(Message::Barrier(barrier))))
})
}
WriteFuture::ReceiveFromUpstream { future, .. } => {
let (opt, stream) = future.await;
must_match!(replace(self, WriteFuture::Empty), WriteFuture::ReceiveFromUpstream { write_state, .. } => {
opt
.ok_or_else(|| anyhow!("end of upstream input").into())
.and_then(|result| result.map(|item| {
(stream, write_state, WriteFutureEvent::UpstreamMessageReceived(item))
}))
})
}
WriteFuture::FlushingChunk { future, .. } => {
let (write_state, result) = future.await;
let result = must_match!(replace(self, WriteFuture::Empty), WriteFuture::FlushingChunk { epoch, start_seq_id, end_seq_id, stream, .. } => {
result.map(|(flush_info, vnode_bitmap)| {
(stream, write_state, WriteFutureEvent::ChunkFlushed(FlushedChunkInfo {
epoch,
start_seq_id,
end_seq_id,
flush_info,
vnode_bitmap,
}))
})
});View on GitHub (pinned to 6469eb736d)