risingwavelabs/risingwave · error

rate limit control channel closed

Error message

rate limit control channel closed

What it means

The rate-limited log store writer's select loop receives on a control channel that delivers `RateLimit` updates; when the sender side is dropped the `recv()` yields `None` and the loop bails with this error. It signals the rate-limit control plane is gone, so the writer can no longer be reconfigured and shuts down.

Solutions

  1. Keep the control channel sender alive for the lifetime of the sink writer
  2. Check why the rate-limit manager task exited (logs, cancellation) and fix its shutdown ordering
  3. If intentional shutdown, stop the log store loop first, or treat `None` as a graceful stop rather than an error

Example fix

// before
let (tx, rx) = mpsc::unbounded_channel();
spawn_sink_writer(rx);
drop(tx); // causes 'rate limit control channel closed'
// after
let (tx, rx) = mpsc::unbounded_channel();
let writer = spawn_sink_writer(rx);
// keep tx alive until writer finishes
writer.await;
drop(tx);
Defensive patterns

Strategy: try-catch

Try / catch

match log_store_result {
    Err(e) if e.to_string().contains("rate limit control channel closed") => {
        info!("rate limit manager gone; stopping sink writer gracefully");
    }
    other => other?,
}

Prevention

When it happens

Trigger: The task/component holding `control_tx` (the rate limit manager) is dropped or terminated while `LogStore` is still running its main loop; `self.control_rx.recv()` returns `None`.

Common situations: Meta node dropping the sink rate-limit control stream; actor shutdown ordering where the manager is cancelled before the sink; network failure tearing down the control gRPC stream.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/log_store.rs:611

                Ok((epoch, LogStoreReadItem::StreamChunk { chunk, chunk_id }))
            }
        }
    }
}

impl<R: LogReader> LogReader for RateLimitedLogReader<R> {
    async fn init(&mut self) -> LogStoreResult<()> {
        self.core.inner.init().await
    }

    async fn next_item(&mut self) -> LogStoreResult<(u64, LogStoreReadItem)> {
        loop {
            select! {
                biased;
                recv = pin!(self.control_rx.recv()) => {
                    let new_rate_limit = match recv {
                        Some(limit) => limit,
                        None => bail!("rate limit control channel closed"),
                    };
                    let old_rate_limit = self.core.rate_limiter.update(new_rate_limit);
                    let paused = matches!(new_rate_limit, RateLimit::Pause);
                    tracing::info!("rate limit changed from {:?} to {:?}, paused = {paused}", old_rate_limit, new_rate_limit);
                },
                item = self.core.next_item() => {
                    return item;
                }
            }
        }
    }

    fn truncate(&mut self, offset: TruncateOffset) -> LogStoreResult<()> {
        let downstream_offset = DownstreamChunkOffset(offset);
        let mut truncate_offset = None;
        let mut stop = false;
        'outer: while let Some((upstream_offset, downstream_offsets)) =
            self.core.consumed_offset_queue.back_mut()

View on GitHub (pinned to 6469eb736d)