{"record":{"id":"51559c9f917626fe","repo":"risingwavelabs/risingwave","slug":"failed-to-inject-offsets-for-splits","errorCode":null,"errorMessage":"failed to inject offsets for splits: {:?}","messagePattern":"failed to inject offsets for splits: (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/source/source_executor.rs","lineNumber":454,"sourceCode":"                        error = ?e.as_report(),\n                        \"Failed to update split offset\"\n                    );\n                    failed_splits.push(split_id.clone());\n                    continue;\n                }\n                // Mark this split as updated for persistence\n                self.stream_source_core\n                    .updated_splits_in_epoch\n                    .insert(split_id.clone().into(), split.clone());\n                // Parse the offset as JSON and store it\n                let json_value: serde_json::Value = serde_json::from_str(offset)\n                    .unwrap_or_else(|_| serde_json::json!({ \"offset\": offset }));\n                json_states.push((split_id.clone(), JsonbVal::from(json_value)));\n            }\n        }\n\n        if !failed_splits.is_empty() {\n            return Err(StreamExecutorError::connector_error(anyhow!(\n                \"failed to inject offsets for splits: {:?}\",\n                failed_splits\n            )));\n        }\n\n        let num_injected = json_states.len();\n        if num_injected > 0 {\n            // Store the injected offsets as JSON in the state table\n            self.stream_source_core\n                .split_state_store\n                .set_states_json(json_states)\n                .await?;\n\n            tracing::info!(\n                actor_id = %self.actor_ctx.id,\n                source_id = %self.stream_source_core.source_id,\n                num_injected = num_injected,\n                \"Offset injection completed for owned splits, triggering rebuild\"","sourceCodeStart":436,"sourceCodeEnd":472,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/source/source_executor.rs#L436-L472","documentation":"When handling a CommandInjectSourceOffsets barrier command, the source executor tries to apply the given offsets to each of its splits via the connector's state machinery. Splits whose offset injection fails are collected into failed_splits and reported wholesale with this connector error.","triggerScenarios":"handle_inject_source_offsets receives inject offsets for split IDs; one or more split_ids are not found in the current split set or the connector rejects the offset format, so those ids land in failed_splits.","commonSituations":"Meta-side CDC/backfill offset injection racing with a split change (split removed before injection), mismatched split_id formatting, or connector state not yet initialized for a newly added split.","solutions":["Check the listed split ids against the source's current split assignments; retry the inject after the split change settles.","Verify the offset format expected by the connector matches what is being injected.","If split IDs were recently changed (rescale/split change), re-issue the command on the new actor layout."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// before injecting offsets, verify split ids exist in the executor's current split set\nlet known: HashSet<_> = executor.split_ids();\nlet unknown: Vec<_> = requested.iter().filter(|id| !known.contains(*id)).collect();\nif !unknown.is_empty() { return Err(skip_or_retry(unknown)); }","typeGuard":null,"tryCatchPattern":"match result {\n    Err(e) if e.to_string().starts_with(\"failed to inject offsets for splits\") => {\n        parse_failed_splits(&e).map_or_else(recover, retry_after_split_settle);\n    }\n    other => propagate(other),\n}","preventionTips":["Serialize split changes and offset injection commands; don't race them","Validate split IDs against the current source layout before injection","Use connector-native offset formats when constructing injection payloads"],"tags":["cdc","split-management","connector"],"backgroundTag":"unexpected-response-shape","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"}