{"record":{"id":"02b410656a0660fb","repo":"risingwavelabs/risingwave","slug":"seek-to-latest-is-not-supported-for-this-connector","errorCode":null,"errorMessage":"seek_to_latest is not supported for this connector","messagePattern":"seek_to_latest is not supported for this connector","errorType":"error_code","errorClass":"ConnectorError","httpStatus":null,"severity":"error","filePath":"src/connector/src/source/base.rs","lineNumber":621,"sourceCode":"        parser_config: ParserConfig,\n        source_ctx: SourceContextRef,\n        columns: Option<Vec<Column>>,\n    ) -> crate::error::ConnectorResult<Self>;\n\n    fn into_stream(self) -> BoxSourceChunkStream;\n\n    fn into_event_stream(self) -> BoxSourceReaderEventStream {\n        self.into_stream()\n            .map_ok(SourceReaderEvent::DataChunk)\n            .boxed()\n    }\n\n    fn backfill_info(&self) -> HashMap<SplitId, BackfillInfo> {\n        HashMap::new()\n    }\n\n    async fn seek_to_latest(&mut self) -> Result<Vec<SplitImpl>> {\n        Err(anyhow!(\"seek_to_latest is not supported for this connector\").into())\n    }\n}\n\n/// Information used to determine whether we should start and finish source backfill.\n///\n/// XXX: if a connector cannot provide the latest offsets (but we want to make it shareable),\n/// perhaps we should ban blocking DDL for it.\n#[derive(Debug, Clone)]\npub enum BackfillInfo {\n    HasDataToBackfill {\n        /// The last available offsets for each split (**inclusive**).\n        ///\n        /// This will be used to determine whether source backfill is finished when\n        /// there are no _new_ messages coming from upstream `SourceExecutor`. Otherwise,\n        /// blocking DDL cannot finish until new messages come.\n        ///\n        /// When there are upstream messages, we will use the latest offsets from the upstream.\n        latest_offset: String,","sourceCodeStart":603,"sourceCodeEnd":639,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/source/base.rs#L603-L639","documentation":"The connector's SplitEnumerator uses the default SourceEnumerator trait implementation of seek_to_latest, which is a stub that always errors. seek_to_latest is only implemented by connectors capable of discovering their latest offsets (e.g. Kafka); others (file, S3, etc.) never support it. It is used when starting a shareable/backfill source from the latest offset.","triggerScenarios":"Calling seek_to_latest on a SplitEnumerator of a connector that did not override the default method; enabling source backfill / shared-source mode on a connector without latest-offset support.","commonSituations":"Creating a shared source (SOURCE_BACKFILL enabled) on a connector like S3/datagen/file where latest offsets cannot be queried; version where a connector lacks the trait override.","solutions":["Use a connector that implements seek_to_latest (e.g. Kafka) for shared/backfill sources","Disable the backfill/shared-source option so the engine does not call seek_to_latest","Implement seek_to_latest for the custom connector's SplitEnumerator"],"exampleFix":"// before (custom connector)\n// (no seek_to_latest override -> default error)\n// after\nasync fn seek_to_latest(&mut self) -> Result<Vec<SplitImpl>> {\n    Ok(vec![SplitImpl::MySplit(MySplit::new_latest())])\n}","handlingStrategy":"try-catch","validationCode":"let supported = matches!(connector, \"kafka\" | \"pulsar\");\nif !supported { /* avoid enabling backfill/seek_to_latest path */ }","typeGuard":null,"tryCatchPattern":"match enumerator.seek_to_latest().await {\n    Err(e) if e.to_string().contains(\"seek_to_latest is not supported\") => {\n        // fall back to non-shareable source creation\n    }\n    other => other?,\n}","preventionTips":["Only enable source backfill/shared mode for connectors known to implement seek_to_latest (Kafka, etc.)","When adding a new connector, implement seek_to_latest instead of relying on the trait default"],"tags":["rust","connector","trait-default","backfill"],"backgroundTag":"method-not-implemented","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"}