Spark流任务Executor内任务耗时不均问题排查求助
Hey there, let's break down your Spark streaming latency issue based on all the details you've shared—this is a tricky but common scenario, so let's walk through it step by step.
Your Problem Context
You're seeing a consistent issue across multiple Spark streaming tasks: one task per Executor takes way longer than others, even though the input volume across tasks is roughly the same. Here's a quick recap of your workflow and cluster setup:
- Source: Direct Kafka stream with 40 partitions
- Transformation steps:
mapParitionsWithPair(flatMap) to expand events into more objects →reduceByKeyaggregation → write results to database - Cluster: Apache Mesos with 2 nodes, 2 cores each
- The latency skew shows up most prominently in the reduce stage timeline
What You've Already Tried (And What We Learned)
You've done some solid debugging already, so let's recap those efforts and the key takeaways:
- Swapped reduce operations: Replaced
reduceByKeywith Java/Kotlin Sequence reduce, but the latency skew persisted—ruling out shuffle-specific issues as the root cause. - Input volume testing:
- High input (160K events): Total runtime 1.8–4.8 mins (~580 events/sec), and the slow task's impact was far smaller than in low-input scenarios.
- Low input: Processing rate fluctuated wildly (660–54 events/sec).
- Interesting note: The slow task's runtime was almost identical (~41s) in both scenarios—this suggests it's not about data size, but a fixed-duration delay.
- Memory adjustments: Added more Executor memory (30% idle now), but the problem didn't go away—so memory pressure isn't the culprit.
- Workflow restructuring:
- Used Java 8 Stream for per-partition reduce to avoid shuffle, adjusted the DAG.
- Increased batch interval to 20s and added more nodes.
- Result: Instead of one slow task, you now have multiple slow tasks + a few fast ones, but overall performance is way better than the short-interval setup.
- Logging & monitoring: Added logs around partition processing, and saw tasks alternate between slow and fast runs. During slow periods, CPU was idle, no network/CPU I/O blocks, and memory stayed steady at 50%.
Key Log Analysis
Let's dig into the logs you shared to spot patterns:
Partition Processing Logs (Only object transformation, no Cassandra write time)
started processing partitioned input: thread 99 started processing partitioned input: thread 98 finished processing partitioned input: thread 99 took 40615ms finished processing partitioned input: thread 98 took 40469ms started processing partitioned input: thread 98 started processing partitioned input: thread 99 finished processing partitioned input: thread 98 took 40476ms finished processing partitioned input: thread 99 took 40523ms started processing partitioned input: thread 98 started processing partitioned input: thread 99 finished processing partitioned input: thread 98 40465ms finished processing partitioned input: thread 99 40379ms started processing partitioned input: thread 98 finished processing partitioned input: thread 98 468 started processing partitioned input: thread 99 finished processing partitioned input: thread 99 525 started processing partitioned input: thread 99 started processing partitioned input: thread 98 finished processing partitioned input: thread 98 738 finished processing partitioned input: thread 99 790 started processing partitioned input: thread 98 finished processing partitioned input: thread 98 took 558 started processing partitioned input: thread 99 finished processing partitioned input: thread 99 took 461 started processing partitioned input: thread 98 finished processing partitioned input: thread 98 took 483 started processing partitioned input: thread 99 finished processing partitioned input: thread 99 took 513 started processing partitioned input: thread 98 finished processing partitioned input: thread 98 took 485 started processing partitioned input: thread 99 finished processing partitioned input: thread 99 took 454
The big red flag here is the periodic alternation between ~40s long runs and sub-1s short runs for threads 98 and 99. This isn't random—it points to a recurring block or resource wait that kicks in at regular intervals.
Cassandra Write Logs (Consistently fast, no latency here)
18/02/07 07:41:47 INFO Executor: Running task 17.0 in stage 5.0 (TID 207) 18/02/07 07:41:47 INFO TorrentBroadcast: Started reading broadcast variable 5 18/02/07 07:41:47 INFO MemoryStore: Block broadcast_5_piece0 stored as bytes in memory (estimated size 7.8 KB, free 1177.1 MB) 18/02/07 07:41:47 INFO TorrentBroadcast: Reading broadcast variable 5 took 33 ms 18/02/07 07:41:47 INFO MemoryStore: Block broadcast_5 stored as values in memory (estimated size 16.4 KB, free 1177.1 MB) 18/02/07 07:41:47 INFO BlockManager: Found block rdd_30_2 locally 18/02/07 07:41:47 INFO BlockManager: Found block rdd_30_17 locally 18/02/07 07:42:02 INFO TableWriter: Wrote 28926 rows to keyspace.table in 15.749 s. 18/02/07 07:42:02 INFO Executor: Finished task 17.0 in stage 5.0 (TID 207). 923 bytes result sent to driver 18/02/07 07:42:02 INFO CoarseGrainedExecutorBackend: Got assigned task 209 18/02/07 07:42:02 INFO Executor: Running task 18.0 in stage 5.0 (TID 209) 18/02/07 07:42:02 INFO BlockManager: Found block rdd_30_18 locally 18/02/07 07:42:03 INFO TableWriter: Wrote 29288 rows to keyspace.table in 16.042 s. 18/02/07 07:42:03 INFO Executor: Finished task 2.0 in stage 5.0 (TID 203). 1713 bytes result sent to driver 18/02/07 07:42:03 INFO CoarseGrainedExecutorBackend: Got assigned task 211 18/02/07 07:42:03 INFO Executor: Running task 21.0 in stage 5.0 (TID 211) 18/02/07 07:42:03 INFO BlockManager: Found block rdd_30_21 locally 18/02/07 07:42:19 INFO TableWriter: Wrote 29315 rows to keyspace.table in 16.308 s. 18/02/07 07:42:19 INFO Executor: Finished task 21.0 in stage 5.0 (TID 211). 923 bytes result sent to driver 18/02/07 07:42:19 INFO CoarseGrainedExecutorBackend: Got assigned task 217 18/02/07 07:42:19 INFO Executor: Running task 24.0 in stage 5.0 (TID 217) 18/02/07 07:42:19 INFO BlockManager: Found block rdd_30_24 locally 18/02/07 07:42:19 INFO TableWriter: Wrote 29422 rows to keyspace.table in 16.783 s. 18/02/07 07:42:19 INFO Executor: Finished task 18.0 in stage 5.0 (TID 209). 923 bytes result sent to driver 18/02/07 07:42:19 INFO CoarseGrainedExecutorBackend: Got assigned task 218 18/02/07 07:42:19 INFO Executor: Running task 25.0 in stage 5.0 (TID 218) 18/02/07 07:42:19 INFO BlockManager: Found block rdd_30_25 locally 18/02/07 07:42:35 INFO TableWriter: Wrote 29427 rows to keyspace.table in 16.509 s. 18/02/07 07:42:35 INFO Executor: Finished task 24.0 in stage 5.0 (TID 217). 923 bytes result sent to driver 18/02/07 07:42:35 INFO CoarseGrainedExecutorBackend: Got assigned task 225
Good news here: Writing to Cassandra is consistently fast and stable, so we can rule out database operations as the source of the latency skew.
Next Steps for Debugging
Based on what we know, here are targeted directions to dig deeper:
- Check Mesos scheduling for periodic interference
Since your cluster runs on Mesos, it's possible that periodic Mesos operations (like resource reclamation, node health checks, or scheduler rebalancing) are pausing Executor threads or stealing resources. Check Mesos master/agent logs for events that line up with the ~40s slow task intervals. - Investigate JVM-level periodic operations
- GC activity: Even though memory usage is stable, a Full GC could be causing pauses. Enable detailed GC logs and see if slow task timestamps match GC events.
- JIT compilation: JVM's just-in-time compiler can cause periodic pauses as it optimizes code. Check JVM logs for compilation events around slow task runs.
- Audit Kafka stream partition consumption
Even if input volume looks similar, check if specific Kafka partitions have delayed messages, large payloads, or malformed data that could cause unexpected processing delays. Also verify the Direct Stream partition assignment strategy to ensure no partitions are being starved or overloaded. - Capture thread dumps during slow runs
When tasks are running slow, usejstackto grab thread dumps of the Executor JVM. Look for threads inWAITINGorBLOCKEDstates—this will pinpoint exactly where the code is stuck (even if you don't see obvious I/O blocks). - Tweak Spark task locality settings
Checkspark.locality.waitand related configuration—if tasks are waiting for local data too long, it could cause unnecessary delays. Also ensure your Executor core count (spark.executor.cores) matches your node's core count (2 per node) to avoid resource contention.
Hopefully these steps help you track down the root cause!
内容的提问来源于stack exchange,提问作者reith

