{"record":{"id":"755da8e71fce993f","repo":"risingwavelabs/risingwave","slug":"only-support-single-split","errorCode":null,"errorMessage":"only support single split","messagePattern":"only support single split","errorType":"validation","errorClass":"ConnectorError","httpStatus":null,"severity":"error","filePath":"src/connector/src/source/pulsar/source/reader.rs","lineNumber":95,"sourceCode":"const PULSAR_DEFAULT_SUBSCRIPTION_PREFIX: &str = \"rw-consumer\";\n\npub enum PulsarSplitReader {\n    Broker(PulsarBrokerReader),\n}\n\n#[async_trait]\nimpl SplitReader for PulsarSplitReader {\n    type Properties = PulsarProperties;\n    type Split = PulsarSplit;\n\n    async fn new(\n        props: PulsarProperties,\n        splits: Vec<PulsarSplit>,\n        parser_config: ParserConfig,\n        source_ctx: SourceContextRef,\n        _columns: Option<Vec<Column>>,\n    ) -> ConnectorResult<Self> {\n        ensure!(splits.len() == 1, \"only support single split\");\n        let split = splits.into_iter().next().unwrap();\n        let topic = split.topic.to_string();\n\n        tracing::debug!(\"creating consumer for pulsar split topic {}\", topic,);\n\n        if props.iceberg_loader_enabled.unwrap_or(false) {\n            bail!(\"PulsarIcebergReader has already been deprecated\");\n        } else {\n            Ok(Self::Broker(\n                PulsarBrokerReader::new(props, vec![split], parser_config, source_ctx, None)\n                    .await?,\n            ))\n        }\n    }\n\n    fn into_stream(self) -> BoxSourceChunkStream {\n        match self {\n            Self::Broker(reader) => {","sourceCodeStart":77,"sourceCodeEnd":113,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/source/pulsar/source/reader.rs#L77-L113","documentation":"The Pulsar source reader supports exactly one split (topic partition) per reader instance; constructing PulsarSplitReader with more than one split is rejected. Pulsar consumers in this connector are created per topic, so multi-split readers are not implemented.","triggerScenarios":"Calling PulsarSplitReader::new (broker reader path) with a splits Vec containing 2 or more PulsarSplits.","commonSituations":"Configuring a partitioned Pulsar topic that produces multiple splits; upstream enumerator changes returning several splits; hand-constructing splits in tests or tooling.","solutions":["Limit the source to a single split/partition, or rely on the framework to distribute splits across parallel reader instances.","If reading a partitioned topic, partition the workload so each reader gets one partition's split.","File a feature request / check newer versions for multi-split Pulsar reader support."],"exampleFix":"// before\nreader.new(props, vec![split_a, split_b], ...).await?;\n// after\nfor split in splits.into_iter() {\n    reader.new(props.clone(), vec![split], ...).await?;\n}","handlingStrategy":"validation","validationCode":"if splits.len() != 1 {\n    return Err(format!(\"pulsar reader accepts exactly 1 split, got {}\", splits.len()));\n}","typeGuard":null,"tryCatchPattern":"match reader.new(props, splits, ...).await {\n    Err(e) if e.to_string().contains(\"only support single split\") => distribute_splits_across_readers(),\n    other => other,\n}","preventionTips":["Assign one Pulsar split per reader instance in your parallelism configuration.","For partitioned topics, ensure split assignment distributes partitions across readers.","Document the single-split constraint in connector integration code/tests."],"tags":["pulsar","source","limitation"],"backgroundTag":"unsupported-operation","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"}