apache/beam · error · IllegalStateException

Spark Receiver was not initialized

Error message

Spark Receiver was not initialized

What it means

ReadFromSparkReceiverWithOffsetDoFn.start() instantiates a WrappedSupervisor to drive a Spark Receiver within Beam. If supervisor construction throws, the DoFn logs the cause and rethrows this IllegalStateException, because the receiver cannot stream records without a running supervisor.

Solutions

  1. Check the logged "Can not init Spark Receiver!" stack trace for the root cause
  2. Ensure spark-core and the receiver's dependencies are bundled in the pipeline jar / available on workers
  3. Provide required SparkConf settings via the receiver's builder options
  4. Verify the SparkReceiverBuilder returns a valid Receiver compatible with sparkreceiver/3 connector

Example fix

// before
SparkReceiverIO.readFromSparkReceiver().withSparkReceiverBuilder(() -> new MyReceiver()) // receiver needs config
// after
SparkReceiverIO.readFromSparkReceiver().withSparkReceiverBuilder(() -> new MyReceiver(sparkConfParams));
Defensive patterns

Strategy: try-catch

Validate before calling

// preflight: receiver must build and Spark deps present
Receiver<?> r = sparkReceiverBuilder.build(); // throws early if broken

Try / catch

try {
  readFromSparkReceiver();
} catch (IllegalStateException e) {
  if (e.getMessage().contains("Spark Receiver was not initialized")) {
    LOG.error("Check SparkConf and spark-core on worker classpath", e);
  }
  throw e;
}

Prevention

When it happens

Trigger: new WrappedSupervisor(sparkReceiver, new SparkConf(), storeFn) throws — e.g. invalid SparkConf, missing Spark dependencies at runtime, receiver constructor failures, or storeFn wiring errors.

Common situations: Missing spark-core dependency in the worker classpath; incompatible Spark Receiver implementation requiring SparkConf settings not provided; receiver builder returning an instance whose constructor fails on the worker environment.

Related errors


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/d7438d88f3541276. Report an issue: GitHub.

Appendix: source

Thrown at sdks/java/io/sparkreceiver/3/src/main/java/org/apache/beam/sdk/io/sparkreceiver/ReadFromSparkReceiverWithOffsetDoFn.java:248

            } else if (data instanceof ArrayBuffer) {
              final ArrayBuffer<V> arrayBuffer = (ArrayBuffer<V>) data;
              final Iterator<V> iterator = arrayBuffer.iterator();
              while (iterator.hasNext()) {
                V record = iterator.next();
                recordsQueue.offer(record);
              }
            } else {
              V record = (V) data;
              recordsQueue.offer(record);
            }
            return null;
          };

      try {
        new WrappedSupervisor(sparkReceiver, new SparkConf(), storeFn);
      } catch (Exception e) {
        LOG.error("Can not init Spark Receiver!", e);
        throw new IllegalStateException("Spark Receiver was not initialized");
      }
      LOG.debug("Starting receiver");
      ((HasOffset) sparkReceiver).setStartOffset(startOffset);
      sparkReceiver.supervisor().startReceiver();
      LOG.debug("Receiver started");
    }

    @Override
    public void stop() {
      if (sparkReceiver != null) {
        sparkReceiver.stop("SparkReceiver is stopped.");
      }
      LOG.info("Clear records queue: {} records", recordsQueue.size());
      recordsQueue.clear();
    }
  }

  @ProcessElement

View on GitHub (pinned to 12126d8942)