risingwavelabs/risingwave · error

invalid persisted stream scan type

Error message

invalid persisted stream scan type {} in job {} fragment {}

What it means

The meta service persisted a stream scan node in a fragment whose stream_scan_type value does not map to any known StreamScanType enum variant in the protobuf definition. This indicates a corrupt or forward-incompatible persisted catalog: the stored proto value is not a valid enum discriminant. It is raised while collecting change-log truncate info for a job's state tables.

Solutions

  1. Check RisingWave version consistency; do not mix meta binary versions against one cluster state.
  2. Rebuild or drop the affected streaming job (job_id from the message) and recreate it.
  3. If this appeared after an upgrade/rollback, restore metadata from backup taken at a compatible version.
  4. Inspect the persisted fragment proto in the metadata store to confirm the invalid discriminant.
Defensive patterns

Strategy: validation

Validate before calling

// before depending on persisted scan info
if StreamScanType::try_from(stream_scan.stream_scan_type).is_err() {
    // skip this job / surface unsupported-value instead of proceeding
}

Try / catch

match get_table_change_log_truncate_info().await {
    Err(e) if e.to_string().contains("invalid persisted stream scan type") => /* version skew: pin meta binary version or recreate job */,
    other => other?,
}

Prevention

When it happens

Trigger: Calling get_table_change_log_truncate_info when a fragment's persisted StreamNode contains a stream_scan_type value not defined in the current StreamScanType proto enum (e.g. meta built from older/newer code reading a catalog written by a different version).

Common situations: Version skew between meta node and persisted metadata (upgrade/rollback), corrupted cluster state after failed migration, or manually edited catalog storage.

Understand the failure class

Background: Invalid enum value errors: "Unknown type", "Invalid scope", "must be one of" — when a string is not on the library's allowed list — this error's family across 23 libraries.

Related errors


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

Appendix: source

Thrown at src/meta/src/controller/streaming_job.rs:311

            })
            .collect();
        if !job_info.is_empty() {
            let fragments = Fragment::find()
                .filter(fragment::Column::JobId.is_in(job_info.keys().copied()))
                .all(&inner.db)
                .await?;
            for fragment in fragments {
                let info = job_info
                    .get_mut(&fragment.job_id)
                    .expect("job should exist");
                info.state_table_ids
                    .extend(fragment.state_table_ids.inner_ref().iter().copied());
                let mut collection_error = None;
                visit_stream_node_stream_scan(&fragment.stream_node.to_protobuf(), |stream_scan| {
                    let scan_type = match StreamScanType::try_from(stream_scan.stream_scan_type) {
                        Ok(scan_type) => scan_type,
                        Err(err) => {
                            collection_error = Some(anyhow::Error::new(err).context(format!(
                                "invalid persisted stream scan type {} in job {} fragment {}",
                                stream_scan.stream_scan_type, fragment.job_id, fragment.fragment_id
                            )));
                            return;
                        }
                    };
                    if scan_type != StreamScanType::SnapshotBackfill {
                        return;
                    }
                    match info
                        .upstream_table_snapshot_epochs
                        .entry(stream_scan.table_id)
                    {
                        std::collections::hash_map::Entry::Occupied(entry) => {
                            if entry.get() != &stream_scan.snapshot_backfill_epoch {
                                collection_error = Some(anyhow!(
                                    "job {} has inconsistent snapshot epochs for upstream table {}",
                                    fragment.job_id,

View on GitHub (pinned to 6469eb736d)