risingwavelabs/risingwave · error

the pubsub reader only supports a single split

Error message

the pubsub reader only supports a single split

What it means

The Google Pub/Sub source reader enforces that it is constructed with exactly one split (PubsubSplit). Pub/Sub subscription reading in RisingWave is inherently a single-stream consumer per reader, so multiple splits cannot be multiplexed by one reader instance. If more than one split is passed, `ensure!` fails and the constructor returns this error.

Solutions

  1. Ensure the upstream split planner/assigner sends at most one PubsubSplit per reader instance.
  2. Verify parallelism configuration: each worker should receive exactly one split for the Pubsub source.
  3. If 0 splits arrive, check that the split enumeration for the subscription produced a split before constructing the reader.
  4. Add a pre-check on the caller side: only call the reader factory when splits.len() == 1.

Example fix

// before: passing whatever splits arrived
let reader = PubsubSplitReader::new(props, all_splits, cfg, ctx, cols).await?;
// after: guarantee a single split per reader
assert!(all_splits.len() <= 1, "pubsub expects one split per reader");
let reader = PubsubSplitReader::new(props, all_splits, cfg, ctx, cols).await?;
Defensive patterns

Strategy: validation

Validate before calling

if splits.len() != 1 {
    return Err(anyhow!("pubsub reader requires exactly one split, got {}", splits.len()));
}
let reader = PubsubSplitReader::new(props, splits, cfg, ctx, cols).await?;

Type guard

fn has_single_split(splits: &[PubsubSplit]) -> bool { splits.len() == 1 }

Prevention

When it happens

Trigger: Constructing `PubsubSplitReader::new` (via the source reader factory) with a `splits: Vec<PubsubSplit>` whose length is not exactly 1 — typically 0 or 2+ splits after split assignment/planning.

Common situations: Running with a source planner that assigns multiple Pub/Sub splits to one worker; misconfigured parallelism where split count != reader count; regression in split-assignment code delivering an empty or multi-element split list to the reader.

Understand the failure class

Background: "Must be a positive integer", "Invalid value", "Unsupported": the invalid-argument-value error family, when a library rejects the value you pass — this error's family across 35 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/source/google_pubsub/source/reader.rs:138

            // Stream ended (all subscribers stopped). Reconnect.
            tracing::warn!("pubsub streaming pull ended, reconnecting...");
        }
    }
}

#[async_trait]
impl SplitReader for PubsubSplitReader {
    type Properties = PubsubProperties;
    type Split = PubsubSplit;

    async fn new(
        properties: PubsubProperties,
        splits: Vec<PubsubSplit>,
        parser_config: ParserConfig,
        source_ctx: SourceContextRef,
        _columns: Option<Vec<Column>>,
    ) -> Result<Self> {
        ensure!(
            splits.len() == 1,
            "the pubsub reader only supports a single split"
        );
        let split = splits.into_iter().next().unwrap();

        let subscriber_config = properties.subscriber_config()?;
        let subscription = properties.subscription_client().await?;

        Ok(Self {
            subscription,
            subscriber_config,
            split_id: split.id(),
            parser_config,
            source_ctx,
        })
    }

    fn into_stream(self) -> BoxSourceChunkStream {

View on GitHub (pinned to 6469eb736d)