{"record":{"id":"b2fa3a8d300cbf69","repo":"risingwavelabs/risingwave","slug":"actor-exited-unexpectedly","errorCode":null,"errorMessage":"actor exited unexpectedly","messagePattern":"actor exited unexpectedly","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/actor.rs","lineNumber":306,"sourceCode":"            ]);\n\n        let mut last_epoch: Option<EpochPair> = None;\n        let mut stream = Box::pin(Box::new(self.consumer).execute());\n\n        // Drive the streaming task with an infinite loop\n        let result = loop {\n            let barrier = match stream\n                .try_next()\n                .instrument(span.clone())\n                .instrument_await(\n                    last_epoch.map_or(await_tree::span!(\"Epoch <initial>\"), |e| {\n                        await_tree::span!(\"Epoch {}\", e.curr)\n                    }),\n                )\n                .await\n            {\n                Ok(Some(barrier)) => barrier,\n                Ok(None) => break Err(anyhow!(\"actor exited unexpectedly\").into()),\n                Err(err) => break Err(err),\n            };\n\n            fail::fail_point!(\"collect_actors_err\", id == 10, |_| Err(anyhow::anyhow!(\n                \"intentional collect_actors_err\"\n            )\n            .into()));\n\n            // Then stop this actor if asked\n            if barrier.is_stop(id) {\n                debug!(actor_id = %id, epoch = ?barrier.epoch, \"stop at barrier\");\n                break Ok(barrier);\n            }\n\n            current_epoch.set(barrier.epoch.curr as i64);\n\n            // Collect barriers to local barrier manager\n            self.barrier_manager.collect(id, &barrier);","sourceCodeStart":288,"sourceCodeEnd":324,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/actor.rs#L288-L324","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["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.","Ensure all upstream actors send a Terminate barrier before their channels are dropped so actors exit via the normal barrier path.","Upgrade RisingWave: this race between channel close and termination has been subject to fixes; use a recent version.","If triggered by tests/fail points (`collect_actors_err`), remove or disable the fail::cfg setup."],"exampleFix":"// before: channel closed silently, actor errors out\nOk(None) => break Err(anyhow!(\"actor exited unexpectedly\").into()),\n// after (upstream): always terminate actors via barrier, never by dropping the sender\n// ensure shutdown path sends Mutation::Stop / Terminate barrier before dropping the barrier tx\nbarrier_tx.send(build_terminate_barrier(prev_epoch)).await?;","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"// In the actor supervising loop:\nmatch actor.run().await {\n    Err(e) if e.to_string().contains(\"actor exited unexpectedly\") => {\n        // treat as abnormal shutdown: alert, then rebuild/restart the streaming job\n        log::warn!(\"actor {} lost its barrier channel; restarting fragment\", actor_id);\n        stream_job_manager.restart_fragment(fragment_id).await?;\n    }\n    other => other?,\n}","preventionTips":["Always terminate actors via Stop/Terminate barriers, never by dropping the barrier sender.","Watch cluster logs for rescheduling/shutdown events overlapping actor errors.","Keep meta and compute nodes on the same version to avoid barrier protocol drift."],"tags":["streaming","actor-lifecycle","barrier","shutdown"],"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"}