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
- Ensure the upstream split planner/assigner sends at most one PubsubSplit per reader instance.
- Verify parallelism configuration: each worker should receive exactly one split for the Pubsub source.
- If 0 splits arrive, check that the split enumeration for the subscription produced a split before constructing the reader.
- 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
- Keep split assignment one-per-reader for pubsub sources
- Validate split count before constructing readers
- Log split count when wiring source readers
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)