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:
Looking at your sink logs and problem description, three key issues are driving the slowdown and eventual Kinesis throughput exception:
Unbounded State Growth in Spark's First Aggregation
Spark 2.1.1 doesn't support automatic state TTL (time-to-live) for streaming aggregations inupdateoutput 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.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.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 theProvisionedThroughputExceededException.
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 useupdatemode, implement a customAggregatorthat 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
Defaultspark.sql.shuffle.partitionsis 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
IncreaseshardFetchIntervalto 2000ms (from default 1000ms) to reduce the frequency of Kinesis read requests, easing throughput pressure:.option("shardFetchInterval", "2000")
After making these changes, monitor these metrics to confirm the fix:
- Spark UI Streaming Tab: Check
State Sizeper batch—this should stabilize instead of growing indefinitely. - Sink Logs: Verify
input rowsandElapsed timeper 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

