{"record":{"id":"d18a5a3d93087f82","repo":"apache/beam","slug":"watermarkestimatorstate-parameters-are-not-supported","errorCode":null,"errorMessage":"@WatermarkEstimatorState parameters are not supported.","messagePattern":"@WatermarkEstimatorState parameters are not supported\\.","errorType":"exception","errorClass":"UnsupportedOperationException","httpStatus":null,"severity":"error","filePath":"sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java","lineNumber":725,"sourceCode":"\n      @Override\n      public @Nullable String currentRecordId(DoFn<InputT, OutputT> doFn) {\n        return outerContext.currentRecordId();\n      }\n\n      @Override\n      public @Nullable Long currentRecordOffset(DoFn<InputT, OutputT> doFn) {\n        return outerContext.currentRecordOffset();\n      }\n\n      @Override\n      public Instant fireTimestamp(DoFn<InputT, OutputT> doFn) {\n        throw new IllegalStateException();\n      }\n\n      @Override\n      public Object watermarkEstimatorState() {\n        throw new UnsupportedOperationException(\n            \"@WatermarkEstimatorState parameters are not supported.\");\n      }\n\n      @Override\n      public WatermarkEstimator<?> watermarkEstimator() {\n        return watermarkEstimator;\n      }\n\n      // ----------- Unsupported methods --------------------\n      @Override\n      public DoFn<InputT, OutputT>.StartBundleContext startBundleContext(\n          DoFn<InputT, OutputT> doFn) {\n        throw new IllegalStateException();\n      }\n\n      @Override\n      public DoFn<InputT, OutputT>.FinishBundleContext finishBundleContext(\n          DoFn<InputT, OutputT> doFn) {","sourceCodeStart":707,"sourceCodeEnd":743,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/SplittableParDoNaiveBounded.java#L707-L743","documentation":"The naive bounded SDF implementation does not expose watermark estimator state to the DoFn. watermarkEstimatorState() on NestedProcessContext always throws, because this fallback path cannot reconstruct per-element watermark estimator state.","triggerScenarios":"A splittable DoFn @ProcessElement method declares a parameter annotated with @WatermarkEstimatorState and is executed via SplittableParDoNaiveBounded.","commonSituations":"Migrating an SDF that tracks watermark estimator state to a runner using the naive bounded implementation; adding @WatermarkEstimatorState parameters for self-checkpointing SDFs on runners that lack full support.","solutions":["Remove the @WatermarkEstimatorState parameter from the DoFn signature","Run the pipeline on a runner/translation with full SDF watermark estimator support","Track estimator state manually inside the DoFn's restriction/checkpoint state instead"],"exampleFix":"// before\n@ProcessElement\npublic void process(ProcessContext c, @WatermarkEstimatorState Instant state) { ... }\n// after\n@ProcessElement\npublic void process(ProcessContext c) { ... } // state not supported here","handlingStrategy":"validation","validationCode":"for (Method m : doFn.getClass().getDeclaredMethods()) {\n  if (m.isAnnotationPresent(ProcessElement.class)) {\n    for (Annotation a : m.getParameterAnnotations()) {\n      if (a instanceof WatermarkEstimatorState) {\n        throw new IllegalStateException(\"@WatermarkEstimatorState unsupported on naive SDF path\");\n      }\n    }\n  }\n}","typeGuard":null,"tryCatchPattern":"try { pipeline.run(); } catch (UnsupportedOperationException e) {\n  if (e.getMessage().contains(\"@WatermarkEstimatorState\")) { /* drop the parameter */ }\n  else throw e;\n}","preventionTips":["Don't declare @WatermarkEstimatorState parameters unless the target runner supports it","Check runner SDF feature matrix before adding watermark estimator parameters"],"tags":["java","beam","sdf","watermark","unsupported"],"backgroundTag":"unsupported-operation","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}