{"record":{"id":"0242614a320e5d5a","repo":"risingwavelabs/risingwave","slug":"the-pubsub-reader-only-supports-a-single-split","errorCode":null,"errorMessage":"the pubsub reader only supports a single split","messagePattern":"the pubsub reader only supports a single split","errorType":"validation","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/connector/src/source/google_pubsub/source/reader.rs","lineNumber":138,"sourceCode":"            // Stream ended (all subscribers stopped). Reconnect.\n            tracing::warn!(\"pubsub streaming pull ended, reconnecting...\");\n        }\n    }\n}\n\n#[async_trait]\nimpl SplitReader for PubsubSplitReader {\n    type Properties = PubsubProperties;\n    type Split = PubsubSplit;\n\n    async fn new(\n        properties: PubsubProperties,\n        splits: Vec<PubsubSplit>,\n        parser_config: ParserConfig,\n        source_ctx: SourceContextRef,\n        _columns: Option<Vec<Column>>,\n    ) -> Result<Self> {\n        ensure!(\n            splits.len() == 1,\n            \"the pubsub reader only supports a single split\"\n        );\n        let split = splits.into_iter().next().unwrap();\n\n        let subscriber_config = properties.subscriber_config()?;\n        let subscription = properties.subscription_client().await?;\n\n        Ok(Self {\n            subscription,\n            subscriber_config,\n            split_id: split.id(),\n            parser_config,\n            source_ctx,\n        })\n    }\n\n    fn into_stream(self) -> BoxSourceChunkStream {","sourceCodeStart":120,"sourceCodeEnd":156,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/source/google_pubsub/source/reader.rs#L120-L156","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// before: passing whatever splits arrived\nlet reader = PubsubSplitReader::new(props, all_splits, cfg, ctx, cols).await?;\n// after: guarantee a single split per reader\nassert!(all_splits.len() <= 1, \"pubsub expects one split per reader\");\nlet reader = PubsubSplitReader::new(props, all_splits, cfg, ctx, cols).await?;","handlingStrategy":"validation","validationCode":"if splits.len() != 1 {\n    return Err(anyhow!(\"pubsub reader requires exactly one split, got {}\", splits.len()));\n}\nlet reader = PubsubSplitReader::new(props, splits, cfg, ctx, cols).await?;","typeGuard":"fn has_single_split(splits: &[PubsubSplit]) -> bool { splits.len() == 1 }","tryCatchPattern":null,"preventionTips":["Keep split assignment one-per-reader for pubsub sources","Validate split count before constructing readers","Log split count when wiring source readers"],"tags":["rust","pubsub","source-reader","split-assignment"],"backgroundTag":"invalid-argument-value","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"}