You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark流任务Executor内任务耗时不均问题排查求助

Troubleshooting Uneven Task Latency in Spark Streaming Executors

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 → reduceByKey aggregation → 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:

  1. Swapped reduce operations: Replaced reduceByKey with Java/Kotlin Sequence reduce, but the latency skew persisted—ruling out shuffle-specific issues as the root cause.
  2. 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.
  3. Memory adjustments: Added more Executor memory (30% idle now), but the problem didn't go away—so memory pressure isn't the culprit.
  4. 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.
  5. 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:

  1. 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.
  2. 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.
  3. 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.
  4. Capture thread dumps during slow runs
    When tasks are running slow, use jstack to grab thread dumps of the Executor JVM. Look for threads in WAITING or BLOCKED states—this will pinpoint exactly where the code is stuck (even if you don't see obvious I/O blocks).
  5. Tweak Spark task locality settings
    Check spark.locality.wait and 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:56:39