Spark Structured Streaming 实战:Kafka 数据源接入与流批一体化处理
·
Spark Structured Streaming 实战:Kafka 数据源接入与流批一体化处理
1. Kafka 数据源接入
Spark Structured Streaming 通过 kafka 格式直接接入 Kafka 数据流,核心配置参数:
kafka.bootstrap.servers: Kafka 集群地址subscribe: 订阅的 TopicstartingOffsets: 消费起始位点(earliest/latest)
接入示例:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("KafkaIntegration") \
.getOrCreate()
# 创建 Kafka 数据流
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:9092,host2:9092") \
.option("subscribe", "iot_topic") \
.option("startingOffsets", "earliest") \
.load()
# 解析 JSON 格式数据(假设 value 为 JSON)
from pyspark.sql.functions import from_json, col
schema = "device_id STRING, temperature DOUBLE, timestamp TIMESTAMP"
parsed_df = df.select(
from_json(col("value").cast("string"), schema).alias("data")
).select("data.*")
2. 流批一体化处理
利用 Spark 统一引擎实现流批统一逻辑:
- 批处理语义:通过微批次(Micro-Batch)实现近实时处理
- 统一 API:
groupBy/window等操作同时适用于流与批 - 结果输出:支持多种 Sink(Kafka/Console/File)
流式聚合示例:
# 按设备ID和10分钟窗口聚合温度
from pyspark.sql.functions import window
agg_df = parsed_df \
.withWatermark("timestamp", "5 minutes") \ # 水位线防迟滞数据
.groupBy(
col("device_id"),
window(col("timestamp"), "10 minutes")
) \
.agg({"temperature": "avg"}) \
.withColumnRenamed("avg(temperature)", "avg_temp")
# 输出到控制台
query = agg_df.writeStream \
.outputMode("update") \
.format("console") \
.option("truncate", "false") \
.start()
query.awaitTermination()
3. 关键配置优化
| 参数 | 推荐值 | 作用 |
|---|---|---|
spark.sql.shuffle.partitions | 200 | 控制 Shuffle 并行度 |
maxOffsetsPerTrigger | 10000 | 单批次最大消费条数 |
failOnDataLoss | false | 容忍数据丢失 |
4. 端到端示例:实时告警系统
# 1. 接入 Kafka 数据流
input_df = spark.readStream.format("kafka")...
# 2. 解析并过滤异常数据
abnormal_df = input_df.filter(col("temperature") > 100)
# 3. 批处理模式:每小时统计异常次数
from pyspark.sql.functions import current_timestamp
batch_df = abnormal_df.withColumn("batch_time", current_timestamp()) \
.groupBy(window("batch_time", "1 hour")) \
.count()
# 4. 双路输出
# 流式输出到 Kafka
abnormal_df.selectExpr("CAST(device_id AS STRING) AS key",
"to_json(struct(*)) AS value") \
.writeStream \
.format("kafka")...
# 批处理结果写入 HDFS
batch_df.writeStream \
.format("parquet") \
.option("path", "/alerts/hourly") \
.trigger(processingTime="1 hour") \
.start()
5. 处理延迟优化公式
流处理延迟由批次间隔 $T_b$ 和数据处理时间 $T_p$ 决定: $$L = T_b + T_p$$ 通过动态调整 maxOffsetsPerTrigger 可平衡吞吐与延迟: $$\Delta O = \frac{\text{目标吞吐量}}{\text{记录大小}} \times T_b$$
更多推荐
所有评论(0)