{"record":{"id":"7e8bdbf6074a78d5","repo":"apache/beam","slug":"spark-receiver-was-not-built","errorCode":null,"errorMessage":"Spark Receiver was not built!","messagePattern":"Spark Receiver was not built!","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/sparkreceiver/3/src/main/java/org/apache/beam/sdk/io/sparkreceiver/ReadFromSparkReceiverWithOffsetDoFn.java","lineNumber":283,"sourceCode":"\n  @ProcessElement\n  public ProcessContinuation processElement(\n      @Element byte[] element,\n      RestrictionTracker<OffsetRange, Long> tracker,\n      WatermarkEstimator<Instant> watermarkEstimator,\n      OutputReceiver<V> receiver) {\n\n    if (tracker.currentRestriction() != null) {\n      LOG.info(\n          \"Start processing element. Restriction = {}\", tracker.currentRestriction().toString());\n    }\n    SparkConsumer<V> sparkConsumer;\n    Receiver<V> sparkReceiver;\n    try {\n      sparkReceiver = sparkReceiverBuilder.build();\n    } catch (Exception e) {\n      LOG.error(\"Can not build Spark Receiver\", e);\n      throw new IllegalStateException(\"Spark Receiver was not built!\");\n    }\n    LOG.debug(\"Restriction {}\", tracker.currentRestriction().toString());\n    sparkConsumer = new SparkConsumerWithOffset<>(tracker.currentRestriction().getFrom());\n    sparkConsumer.start(sparkReceiver);\n\n    Long recordsProcessed = 0L;\n    while (true) {\n      LOG.debug(\"Start polling records\");\n      try {\n        TimeUnit.SECONDS.sleep(startPollTimeoutSec);\n      } catch (InterruptedException e) {\n        LOG.error(\"SparkReceiver was interrupted before polling started\", e);\n        throw new IllegalStateException(\"Spark Receiver was interrupted before polling started\");\n      }\n      if (!sparkConsumer.hasRecords()) {\n        LOG.debug(\"No records left\");\n        ((HasOffset) sparkReceiver).setCheckpoint(recordsProcessed);\n        sparkConsumer.stop();","sourceCodeStart":265,"sourceCodeEnd":301,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/sparkreceiver/3/src/main/java/org/apache/beam/sdk/io/sparkreceiver/ReadFromSparkReceiverWithOffsetDoFn.java#L265-L301","documentation":"In processElement, the DoFn builds a fresh Spark Receiver instance per element via the user-supplied SparkReceiverBuilder. If build() throws, the receiver cannot be created and this IllegalStateException aborts processing of the element.","triggerScenarios":"sparkReceiverBuilder.build() throws for the given element — user-supplied lambda fails (bad constructor args, missing external resources, unsupported receiver state), or the receiver requires context not present on the Beam worker.","commonSituations":"Builder lambdas capturing non-serializable or worker-unavailable state; receivers that open network connections in their constructor and fail due to environment/firewall; bugs in custom Receiver constructors.","solutions":["Inspect the logged \"Can not build Spark Receiver\" stack trace for the root cause","Make the builder lambda static/stateless and serializable, sourcing any config from pipeline options","Ensure the Receiver's constructor has no worker-environment-dependent failures (ports, files, credentials)","Test the receiver's no-arg construction outside Beam to confirm it builds cleanly"],"exampleFix":"// before\n.withSparkReceiverBuilder(() -> new MyReceiver(nonSerializableHelper.getConfig()))\n// after\n.withSparkReceiverBuilder(() -> new MyReceiver(options.getReceiverParam()))","handlingStrategy":"validation","validationCode":"// preflight outside pipeline\nReceiver<V> r = sparkReceiverBuilder.build(); // must not throw","typeGuard":"boolean buildsCleanly(Supplier<Receiver<V>> b) { try { return b.get() != null; } catch (Exception e) { return false; } }","tryCatchPattern":"try {\n  processElement();\n} catch (IllegalStateException e) {\n  if (e.getMessage().contains(\"Spark Receiver was not built\")) {\n    LOG.error(\"Receiver builder failed on worker; check serialization and env\", e);\n  }\n  throw e;\n}","preventionTips":["Keep builder lambdas static, stateless and serializable","Pass config via pipeline options, not captured objects","Avoid constructor-time side effects (network/files) in Receivers","Unit-test the receiver's no-arg construction"],"tags":["java","beam","spark","receiver","builder"],"backgroundTag":"invalid-argument-value","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"}