risingwavelabs/risingwave · error

actor exited unexpectedly

Error message

actor exited unexpectedly

What it means

A stream actor's main loop in `Actor::run_consumer` polls the barrier channel via `collect_barrier_from_channel`. When the channel returns `Ok(None)` the channel is closed without an error, meaning no more barriers (including the final `Terminate` barrier) will ever arrive. The actor treats this as an abnormal termination and breaks out with this error instead of exiting cleanly.

Solutions

  1. Check meta/compute logs for cancellation or shutdown events that occurred at the same time as this error; restart the affected actor/fragment by rebuilding the streaming job if it is stuck.
  2. Ensure all upstream actors send a Terminate barrier before their channels are dropped so actors exit via the normal barrier path.
  3. Upgrade RisingWave: this race between channel close and termination has been subject to fixes; use a recent version.
  4. If triggered by tests/fail points (`collect_actors_err`), remove or disable the fail::cfg setup.

Example fix

// before: channel closed silently, actor errors out
Ok(None) => break Err(anyhow!("actor exited unexpectedly").into()),
// after (upstream): always terminate actors via barrier, never by dropping the sender
// ensure shutdown path sends Mutation::Stop / Terminate barrier before dropping the barrier tx
barrier_tx.send(build_terminate_barrier(prev_epoch)).await?;
Defensive patterns

Strategy: try-catch

Try / catch

// In the actor supervising loop:
match actor.run().await {
    Err(e) if e.to_string().contains("actor exited unexpectedly") => {
        // treat as abnormal shutdown: alert, then rebuild/restart the streaming job
        log::warn!("actor {} lost its barrier channel; restarting fragment", actor_id);
        stream_job_manager.restart_fragment(fragment_id).await?;
    }
    other => other?,
}

Prevention

When it happens

Trigger: The upstream barrier sender (stream manager / compute actor context) is dropped or closed while the actor is still running, e.g. the fragment's producers were cancelled or the actor's barrier input channel was closed before a terminate barrier was delivered.

Common situations: Compute node shutdown or `ChangeConfig`/migrate racing with actor execution; a dropped upstream fragment during DDL; internal scheduler cancellation that closes barrier channels without sending terminate barriers; also induced in tests via the `collect_actors_err` fail point.

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


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

Appendix: source

Thrown at src/stream/src/executor/actor.rs:306

            ]);

        let mut last_epoch: Option<EpochPair> = None;
        let mut stream = Box::pin(Box::new(self.consumer).execute());

        // Drive the streaming task with an infinite loop
        let result = loop {
            let barrier = match stream
                .try_next()
                .instrument(span.clone())
                .instrument_await(
                    last_epoch.map_or(await_tree::span!("Epoch <initial>"), |e| {
                        await_tree::span!("Epoch {}", e.curr)
                    }),
                )
                .await
            {
                Ok(Some(barrier)) => barrier,
                Ok(None) => break Err(anyhow!("actor exited unexpectedly").into()),
                Err(err) => break Err(err),
            };

            fail::fail_point!("collect_actors_err", id == 10, |_| Err(anyhow::anyhow!(
                "intentional collect_actors_err"
            )
            .into()));

            // Then stop this actor if asked
            if barrier.is_stop(id) {
                debug!(actor_id = %id, epoch = ?barrier.epoch, "stop at barrier");
                break Ok(barrier);
            }

            current_epoch.set(barrier.epoch.curr as i64);

            // Collect barriers to local barrier manager
            self.barrier_manager.collect(id, &barrier);

View on GitHub (pinned to 6469eb736d)