risingwavelabs/risingwave · error

failed to inject offsets for splits

Error message

failed to inject offsets for splits: {:?}

What it means

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.

Solutions

  1. Check the listed split ids against the source's current split assignments; retry the inject after the split change settles.
  2. Verify the offset format expected by the connector matches what is being injected.
  3. If split IDs were recently changed (rescale/split change), re-issue the command on the new actor layout.
Defensive patterns

Strategy: validation

Validate before calling

// before injecting offsets, verify split ids exist in the executor's current split set
let known: HashSet<_> = executor.split_ids();
let unknown: Vec<_> = requested.iter().filter(|id| !known.contains(*id)).collect();
if !unknown.is_empty() { return Err(skip_or_retry(unknown)); }

Try / catch

match result {
    Err(e) if e.to_string().starts_with("failed to inject offsets for splits") => {
        parse_failed_splits(&e).map_or_else(recover, retry_after_split_settle);
    }
    other => propagate(other),
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/source/source_executor.rs:454

                        error = ?e.as_report(),
                        "Failed to update split offset"
                    );
                    failed_splits.push(split_id.clone());
                    continue;
                }
                // Mark this split as updated for persistence
                self.stream_source_core
                    .updated_splits_in_epoch
                    .insert(split_id.clone().into(), split.clone());
                // Parse the offset as JSON and store it
                let json_value: serde_json::Value = serde_json::from_str(offset)
                    .unwrap_or_else(|_| serde_json::json!({ "offset": offset }));
                json_states.push((split_id.clone(), JsonbVal::from(json_value)));
            }
        }

        if !failed_splits.is_empty() {
            return Err(StreamExecutorError::connector_error(anyhow!(
                "failed to inject offsets for splits: {:?}",
                failed_splits
            )));
        }

        let num_injected = json_states.len();
        if num_injected > 0 {
            // Store the injected offsets as JSON in the state table
            self.stream_source_core
                .split_state_store
                .set_states_json(json_states)
                .await?;

            tracing::info!(
                actor_id = %self.actor_ctx.id,
                source_id = %self.stream_source_core.source_id,
                num_injected = num_injected,
                "Offset injection completed for owned splits, triggering rebuild"

View on GitHub (pinned to 6469eb736d)