{"record":{"id":"20b2f8690233612d","repo":"risingwavelabs/risingwave","slug":"failed-to-build-data-generation-runtime","errorCode":null,"errorMessage":"failed to build data-generation runtime","messagePattern":"failed to build data-generation runtime","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"critical","filePath":"src/connector/src/source/data_gen_util.rs","lineNumber":37,"sourceCode":"use tokio::sync::mpsc;\n\n/// Spawn the data generator to a dedicated runtime, returns a channel receiver\n/// for acquiring the generated data. This is used for the [`DatagenSplitReader`]\n/// and [`NexmarkSplitReader`] in case that they are CPU intensive\n/// and may block the streaming actors.\n///\n/// [`DatagenSplitReader`]: super::datagen::DatagenSplitReader\n/// [`NexmarkSplitReader`]: super::nexmark::source::reader::NexmarkSplitReader\npub fn spawn_data_generation_stream<T: Send + 'static>(\n    stream: impl Stream<Item = T> + Send + 'static,\n    buffer_size: usize,\n) -> impl Stream<Item = T> + Send + 'static {\n    static RUNTIME: LazyLock<Runtime> = LazyLock::new(|| {\n        tokio::runtime::Builder::new_multi_thread()\n            .thread_name(\"rw-datagen\")\n            .enable_all()\n            .build()\n            .expect(\"failed to build data-generation runtime\")\n    });\n\n    let (generation_tx, generation_rx) = mpsc::channel(buffer_size);\n    RUNTIME.spawn(async move {\n        pin_mut!(stream);\n        while let Some(result) = stream.next().await {\n            if generation_tx.send(result).await.is_err() {\n                tracing::warn!(\"failed to send next event to reader, exit\");\n                break;\n            }\n        }\n    });\n\n    tokio_stream::wrappers::ReceiverStream::new(generation_rx)\n}\n","sourceCodeStart":19,"sourceCodeEnd":53,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/source/data_gen_util.rs#L19-L53","documentation":"`spawn_data_generation_stream` lazily builds a dedicated multi-thread Tokio runtime (named \"rw-datagen\") for data-generation streams. If `build()` fails (thread creation failure, resource exhaustion, platform constraints) it panics with this expect message.","triggerScenarios":"First call into any datagen source's `into_stream`/`into_data_stream` when the process cannot spawn a new multi-threaded runtime: OS thread/thread-limit exhaustion, low memory, restricted sandbox environments.","commonSituations":"Containers with very low thread/pid limits (ulimit -u, pids cgroup limit), heavily loaded hosts, or running under test harnesses that restrict threading.","solutions":["Raise the container/OS thread limit (ulimit -u, cgroups pids.max) and free memory, then retry.","Reuse an existing runtime instead of building a dedicated one if embedding in your own Tokio app.","Patch to `handle` an existing runtime or fall back to a current-thread runtime builder when multi-thread build fails."],"exampleFix":"// before\n.build()\n.expect(\"failed to build data-generation runtime\")\n// after\n.build()\n.unwrap_or_else(|e| panic!(\"failed to build data-generation runtime: {e}\")) // with logging, or fallback to current_thread builder","handlingStrategy":"try-catch","validationCode":"// Check thread limits before launching datagen workloads:\nulimit -u   # and cgroup pids.max; must allow dozens of threads","typeGuard":null,"tryCatchPattern":"std::panic::catch_unwind(|| {\n    spawn_data_generation_stream(stream, buffer, gen_fn)\n}).unwrap_or_else(|_| {\n    tracing::error!(\"datagen runtime failed to build; check thread/memory limits\");\n    // fall back to a current-thread runtime or skip datagen\n});","preventionTips":["Raise container thread/pids limits for RisingWave deployments.","Ensure sufficient memory headroom for extra Tokio runtimes.","Test datagen sources in constrained sandboxes before production."],"tags":["datagen","tokio","runtime","resource-exhaustion"],"backgroundTag":"module-init-failed","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"}