## 简介
Spark Streaming 是 Apache Spark 生态中用于处理实时数据流的组件。它可以将流数据按时间片切分成小批次,然后用 Spark 引擎进行处理。
## 核心概念
### DStream
DStream(Discretized Stream)是 Spark Streaming 的核心抽象,代表连续的数据流。它由一系列连续的 RDD 组成。
“`python
from pyspark.streaming import StreamingContext
ssc = StreamingContext(sc, batchDuration=1) # 1秒批次
lines = ssc.socketTextStream(“localhost”, 9999)
counts = lines.flatMap(lambda line: line.split(” “)) \
.map(lambda word: (word, 1)) \
.reduceByKey(lambda a, b: a + b)
counts.pprint()
ssc.start()
ssc.awaitTermination()
“`
## 性能优化
1. **合理设置批次间隔**:通常设置在 500ms 到 5s 之间
2. **使用 Checkpoint**:确保故障恢复能力
3. **调整并行度**:根据集群资源合理分配
4. **反压机制**:开启 `spark.streaming.backpressure.enabled`
## 总结
Spark Streaming 适合秒级延迟的批处理场景,如果对延迟要求极高,可以考虑 Flink 或 Kafka Streams。