{"record":{"id":"2084bbb47ae56bf7","repo":"apache/beam","slug":"given-spark-receiver-class-s-doesn-t-implement-hasoffset","errorCode":null,"errorMessage":"Given Spark Receiver class %s doesn't implement HasOffset interface, therefore it is not supported!","messagePattern":"Given Spark Receiver class (.+?) doesn't implement HasOffset interface, therefore it is not supported!","errorType":"exception","errorClass":"UnsupportedOperationException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/sparkreceiver/3/src/main/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIO.java","lineNumber":187,"sourceCode":"      checkStateNotNull(getGetOffsetFn(), \"withGetOffsetFn() is required\");\n    }\n  }\n\n  static class ReadFromSparkReceiverViaSdf<V> extends PTransform<PBegin, PCollection<V>> {\n\n    private final Read<V> sparkReceiverRead;\n\n    ReadFromSparkReceiverViaSdf(Read<V> sparkReceiverRead) {\n      this.sparkReceiverRead = sparkReceiverRead;\n    }\n\n    @Override\n    public PCollection<V> expand(PBegin input) {\n      final ReceiverBuilder<V, ? extends Receiver<V>> sparkReceiverBuilder =\n          sparkReceiverRead.getSparkReceiverBuilder();\n      checkStateNotNull(sparkReceiverBuilder, \"withSparkReceiverBuilder() is required\");\n      if (!HasOffset.class.isAssignableFrom(sparkReceiverBuilder.getSparkReceiverClass())) {\n        throw new UnsupportedOperationException(\n            String.format(\n                \"Given Spark Receiver class %s doesn't implement HasOffset interface,\"\n                    + \" therefore it is not supported!\",\n                sparkReceiverBuilder.getSparkReceiverClass().getName()));\n      } else {\n        LOG.info(\"{} started reading\", ReadFromSparkReceiverWithOffsetDoFn.class.getSimpleName());\n        return input\n            .apply(Impulse.create())\n            .apply(ParDo.of(new ReadFromSparkReceiverWithOffsetDoFn<>(sparkReceiverRead)));\n        // TODO: Split data from SparkReceiver into multiple workers\n      }\n    }\n  }\n}\n","sourceCodeStart":169,"sourceCodeEnd":202,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/sparkreceiver/3/src/main/java/org/apache/beam/sdk/io/sparkreceiver/SparkReceiverIO.java#L169-L202","documentation":"SparkReceiverIO's Read with offsets mode only supports Spark Receiver classes that implement the HasOffset interface, because offset tracking for restartability requires the receiver to expose its current offset. During pipeline expansion (expand), the transform checks the receiver class supplied via withSparkReceiverBuilder() and refuses to run an incompatible receiver rather than fail later at runtime.","triggerScenarios":"Calling SparkReceiverIO.<T>readWithOffsets().withSparkReceiverBuilder(new ReceiverBuilder<>(SomeReceiver.class)) where SomeReceiver (or a subclass of it) does not implement org.apache.beam.sdk.io.sparkreceiver.HasOffset; the exception is thrown when the pipeline graph is constructed/expanded.","commonSituations":"Using a custom Spark Receiver written for the plain SparkReceiverIO.read() (CustomReceiverWithOffset-free) path and then switching to readWithOffsets(); copying an example receiver from older Beam versions before HasOffset was introduced; forgetting to add offset reporting methods when upgrading the pipeline.","solutions":["Make the receiver class implement the HasOffset interface (add getCurrentOffset returning the receiver's current offset and update it as records arrive).","If offset tracking is genuinely not possible, use SparkReceiverIO.read() (streaming without offsets) instead of readWithOffsets().","Verify the class passed to ReceiverBuilder is the concrete receiver that implements HasOffset, not a wrapper or a wrong generic parameter."],"exampleFix":"// before\nreturn io.apply(SparkReceiverIO.<Long>readWithOffsets()\n    .withSparkReceiverBuilder(new ReceiverBuilder<>(MyReceiver.class)));\n// after\npublic class MyReceiver extends Receiver<Long> implements HasOffset {\n  private long currentOffset = 0;\n  @Override\n  public long getCurrentOffset() { return currentOffset; }\n  // ... update currentOffset in onStart/onStop/receive\n}","handlingStrategy":"validation","validationCode":"// Java — before building the read\nClass<? extends Receiver<Long>> rc = myReceiverClass;\nif (!HasOffset.class.isAssignableFrom(rc)) {\n  throw new IllegalStateException(rc.getName() + \" must implement HasOffset for readWithOffsets()\");\n}","typeGuard":"boolean supportsOffsets(Class<? extends Receiver<?>> c) { return HasOffset.class.isAssignableFrom(c); }","tryCatchPattern":null,"preventionTips":["Implement HasOffset on every custom receiver used with readWithOffsets().","Keep receiver classes for offset mode and plain mode clearly named/separated.","Add a unit test that asserts HasOffset.class.isAssignableFrom(yourReceiver.class)."],"tags":["apache-beam","java","spark-receiver","pipeline-construction"],"backgroundTag":"incompatible-source-type","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}