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

Spark Structured Streaming从Kinesis到DynamoDB延迟异常排查问询

Alright, let's dig into your Spark Structured Streaming issue with Kinesis and DynamoDB—this looks like a classic state management + backpressure problem that's common in older Spark versions. Here's a breakdown of what's going wrong and how to fix it:

Root Cause Analysis

Looking at your sink logs and problem description, three key issues are driving the slowdown and eventual Kinesis throughput exception:

  1. Unbounded State Growth in Spark's First Aggregation
    Spark 2.1.1 doesn't support automatic state TTL (time-to-live) for streaming aggregations in update output mode. As your job runs, Spark keeps accumulating state for every aggregation key it's ever seen. Over time, each microbatch has to process more and more state data, slowing down batch processing. This creates backpressure: Kinesis data keeps piling up, and when Spark finally tries to catch up, it pulls a massive chunk of events at once—exceeding your Kinesis shard throughput limits.

  2. Sink-Side Secondary Aggregation Amplifying the Problem
    Moving aggregation logic to your DynamoDB Sink makes things worse. Since Spark's state isn't being cleaned, each batch sends more and more update records to the sink (including repeated updates to old keys). Your sink then has to process an ever-growing number of rows per batch, increasing latency and further delaying Spark's ability to process new Kinesis data.

  3. Missing Trigger Configuration
    Running without a time trigger means Spark tries to process batches as fast as possible. When processing slows down, there's no cap on how much data Spark tries to pull from Kinesis in one go, which directly leads to the ProvisionedThroughputExceededException.


Solutions & Tuning Steps

Let's fix this step by step, prioritizing the most impactful changes first:

1. Fix Unbounded State Growth (Critical)

Since Spark 2.1.1 lacks built-in state TTL, you have two practical options:

  • Switch to Watermark + Append Mode (If Business Logic Allows)
    If you can adjust your aggregation to use event time, add a watermark to automatically clean up old state:

    val aggregatedStream = rawKinesisStream
      .withWatermark("event_time", "1 hour") // Adjust window to match your data's valid retention age
      .groupBy("key")
      .agg(sum("value").as("total"))
      .writeStream
      .outputMode("append")
    

    This tells Spark to drop state for events older than the watermark, keeping state size manageable.

  • Custom Aggregator with Manual State Cleaning
    If you must use update mode, implement a custom Aggregator that filters out expired state entries during updates. For example, track the last update time for each key and discard entries that are past your retention window.

2. Move Secondary Aggregation to Spark

You mentioned Spark "doesn't support multi-aggregation"—that's not accurate for your use case. You can chain aggregations directly in your streaming query instead of pushing logic to the sink. For example:

// First aggregation (as before)
val firstAgg = rawKinesisStream.groupBy("key1").agg(count("*").as("count1"))
// Second aggregation (moved from sink)
val finalAgg = firstAgg.groupBy("key2").agg(sum("count1").as("total_count"))
// Write to DynamoDB Sink with only final results
finalAgg.writeStream
  .outputMode("update")
  .foreach(dynamoDBSink)
  .start()

This reduces the number of rows sent to your sink, cutting down on processing time and insert requests.

3. Add Trigger & Batch Size Controls

  • Set a Processing Time Trigger
    Force Spark to process batches at fixed intervals to avoid unbounded batch sizes:
    .writeStream
    .trigger(Trigger.ProcessingTime("1 minute")) // Adjust interval based on your data rate
    
  • Limit Kinesis Offsets per Trigger
    Cap how much data Spark reads per batch to stay within Kinesis throughput limits:
    spark.readStream
      .format("kinesis")
      .option("maxOffsetsPerTrigger", "5000") // Adjust based on your shard throughput (1MB/shard/sec)
      // Other Kinesis configs
      .load()
    

4. Tune Kinesis & Spark Resources

  • Match Kinesis Shards to Spark Parallelism
    Your m4.4xlarge has 16 cores—set Kinesis shards to 16 (1:1 with Spark tasks) to maximize parallel processing.
  • Adjust Spark Shuffle Partitions
    Default spark.sql.shuffle.partitions is 200, which is overkill for a 16-core cluster. Set it to 32 (2x core count) to reduce shuffle overhead:
    spark.conf.set("spark.sql.shuffle.partitions", "32")
    
  • Tune Kinesis Fetch Interval
    Increase shardFetchInterval to 2000ms (from default 1000ms) to reduce the frequency of Kinesis read requests, easing throughput pressure:
    .option("shardFetchInterval", "2000")
    

Verification Steps

After making these changes, monitor these metrics to confirm the fix:

  • Spark UI Streaming Tab: Check State Size per batch—this should stabilize instead of growing indefinitely.
  • Sink Logs: Verify input rows and Elapsed time per batch stay consistent (no more upward trend).
  • Kinesis CloudWatch Metrics: Track GetRecords.IteratorAgeMilliseconds—it should stay low and stable, indicating no backlog.

内容的提问来源于stack exchange,提问作者philippe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:37:12