Describe the bug
connectedComponents can return incorrect component assignments with no error when one of its internal eager checkpoints fails to durably persist to HDFS. We observed an HDFS DataStreamer failure (0 datanode(s) running) during the checkpoint write, but the Spark task still reported success (Finished task ... result sent to driver). The algorithm then continued on the (empty/incomplete) checkpoint and produced wrong clusters — edges were silently dropped — without throwing anything to the caller.
This is distinct from #453 (AQE) and #760 (non-deterministic input): in our case there is an explicit HDFS error in the logs, yet the job neither retries nor fails.
To Reproduce
Steps to reproduce the behavior:
sc.setCheckpointDir().
Build a graph large/edgy enough that connectedComponents performs at least one internal checkpoint (≥ checkpointInterval iterations).
Make the HDFS checkpoint target fail to persist during the run while still allowing the executor task to "complete" (we hit this with datanodes unavailable / 0 datanode(s) running transient conditions on a standalone cluster writing to HDFS).
Observe: DataStreamer Exception in executor logs, tasks report Finished, run() returns, and the resulting component column is wrong (edges ignored), with no exception.
(We acknowledge a fully deterministic repro requires inducing a checkpoint-persist failure that still lets the task succeed; the executor log above is the key evidence.)
Expected behavior
If an intermediate checkpoint does not durably persist, connectedComponents should fail (raise an exception) rather than silently return incorrect components. Silent wrong results for an identity/graph algorithm are dangerous.
System [please complete the following information]:
GraphFrames: io.graphframes:graphframes-spark4_2.13:0.10.0
Spark: 4.0.0
Scala: 2.13.16, Java: 17
Deployment: Spark standalone cluster (also reproduced as a hard failure in local mode)
Checkpointing: reliable checkpoint dir on HDFS (Hadoop 3.4.x), default algorithm (two-phase), default broadcastThreshold, checkpointInterval = 2, useLocalCheckpoints = false
Component
Additional context
The checkpoint-write tasks for rdd-1066 hit 0 datanodes, but the same task IDs then finished successfully. Note the temp-file attempt-NNN numbers match the task TIDs:
WARN DataStreamer: DataStreamer Exception
org.apache.hadoop.ipc.RemoteException(java.io.IOException): File /.../spark-checkpoint//rdd-1066/.part-00000-attempt-760
could only be written to 0 of the 1 minReplication nodes. There are 0 datanode(s) running ...
at org.apache.hadoop.hdfs.DataStreamer.run(DataStreamer.java:752)
WARN DataStreamer: DataStreamer Exception ... rdd-1066/.part-00002-attempt-761 ... 0 datanode(s) running ...
WARN DataStreamer: DataStreamer Exception ... rdd-1066/.part-00005-attempt-762 ... 0 datanode(s) running ...
INFO Executor: Finished task 0.0 in stage 559.0 (TID 760). 4944 bytes result sent to driver
INFO Executor: Finished task 2.0 in stage 559.0 (TID 761). 4944 bytes result sent to driver
INFO Executor: Finished task 5.0 in stage 559.0 (TID 762). 4944 bytes result sent to driver
So the HDFS block write failed (logged at WARN by the client's background DataStreamer thread), but the executor tasks reported success. Spark therefore considered the checkpoint stage successful, run() returned normally, and connectedComponents produced wrong components (records that should share a component were split apart).
For contrast, when the entire run is against a permanently-unavailable HDFS (e.g. local mode, 0 datanodes for the whole run), all task attempts fail and the job aborts at the same line — so the difference between "silent wrong" and "loud failure" comes down to whether the checkpoint task reports success.
Where it happens
The failing checkpoint is the eager reliable checkpoint inside the iteration loop:
// ConnectedComponents.scala (0.10.0), ~line 321-327
if (shouldCheckpoint && (iteration % checkpointInterval == 0)) {
if (useLocalCheckpoints) {
ee = ee.localCheckpoint(eager = true)
} else {
ee = ee.checkpoint(eager = true) // <-- line 325: our failure originates here
}
}
Our driver stack confirms the origin is line 325 → Dataset.checkpoint → ReliableCheckpointRDD.writeRDDToCheckpointDirectory → addBlock → 0 datanodes.
Are you planning on creating a PR?
Describe the bug
connectedComponents can return incorrect component assignments with no error when one of its internal eager checkpoints fails to durably persist to HDFS. We observed an HDFS DataStreamer failure (0 datanode(s) running) during the checkpoint write, but the Spark task still reported success (Finished task ... result sent to driver). The algorithm then continued on the (empty/incomplete) checkpoint and produced wrong clusters — edges were silently dropped — without throwing anything to the caller.
This is distinct from #453 (AQE) and #760 (non-deterministic input): in our case there is an explicit HDFS error in the logs, yet the job neither retries nor fails.
To Reproduce
Steps to reproduce the behavior:
sc.setCheckpointDir().
Build a graph large/edgy enough that connectedComponents performs at least one internal checkpoint (≥ checkpointInterval iterations).
Make the HDFS checkpoint target fail to persist during the run while still allowing the executor task to "complete" (we hit this with datanodes unavailable / 0 datanode(s) running transient conditions on a standalone cluster writing to HDFS).
Observe: DataStreamer Exception in executor logs, tasks report Finished, run() returns, and the resulting component column is wrong (edges ignored), with no exception.
(We acknowledge a fully deterministic repro requires inducing a checkpoint-persist failure that still lets the task succeed; the executor log above is the key evidence.)
Expected behavior
If an intermediate checkpoint does not durably persist, connectedComponents should fail (raise an exception) rather than silently return incorrect components. Silent wrong results for an identity/graph algorithm are dangerous.
System [please complete the following information]:
GraphFrames: io.graphframes:graphframes-spark4_2.13:0.10.0
Spark: 4.0.0
Scala: 2.13.16, Java: 17
Deployment: Spark standalone cluster (also reproduced as a hard failure in local mode)
Checkpointing: reliable checkpoint dir on HDFS (Hadoop 3.4.x), default algorithm (two-phase), default broadcastThreshold, checkpointInterval = 2, useLocalCheckpoints = false
Component
Additional context
The checkpoint-write tasks for rdd-1066 hit 0 datanodes, but the same task IDs then finished successfully. Note the temp-file attempt-NNN numbers match the task TIDs:
WARN DataStreamer: DataStreamer Exception
org.apache.hadoop.ipc.RemoteException(java.io.IOException): File /.../spark-checkpoint//rdd-1066/.part-00000-attempt-760
could only be written to 0 of the 1 minReplication nodes. There are 0 datanode(s) running ...
at org.apache.hadoop.hdfs.DataStreamer.run(DataStreamer.java:752)
WARN DataStreamer: DataStreamer Exception ... rdd-1066/.part-00002-attempt-761 ... 0 datanode(s) running ...
WARN DataStreamer: DataStreamer Exception ... rdd-1066/.part-00005-attempt-762 ... 0 datanode(s) running ...
INFO Executor: Finished task 0.0 in stage 559.0 (TID 760). 4944 bytes result sent to driver
INFO Executor: Finished task 2.0 in stage 559.0 (TID 761). 4944 bytes result sent to driver
INFO Executor: Finished task 5.0 in stage 559.0 (TID 762). 4944 bytes result sent to driver
So the HDFS block write failed (logged at WARN by the client's background DataStreamer thread), but the executor tasks reported success. Spark therefore considered the checkpoint stage successful, run() returned normally, and connectedComponents produced wrong components (records that should share a component were split apart).
For contrast, when the entire run is against a permanently-unavailable HDFS (e.g. local mode, 0 datanodes for the whole run), all task attempts fail and the job aborts at the same line — so the difference between "silent wrong" and "loud failure" comes down to whether the checkpoint task reports success.
Where it happens
The failing checkpoint is the eager reliable checkpoint inside the iteration loop:
// ConnectedComponents.scala (0.10.0), ~line 321-327
if (shouldCheckpoint && (iteration % checkpointInterval == 0)) {
if (useLocalCheckpoints) {
ee = ee.localCheckpoint(eager = true)
} else {
ee = ee.checkpoint(eager = true) // <-- line 325: our failure originates here
}
}
Our driver stack confirms the origin is line 325 → Dataset.checkpoint → ReliableCheckpointRDD.writeRDDToCheckpointDirectory → addBlock → 0 datanodes.
Are you planning on creating a PR?