多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

智能家居数据管道实战:Kafka + Spark 流批一体处理

智能家居数据管道实战:Kafka + Spark 流批一体处理 简介这是一套面向物联网与大数据方向学习者的智能家居数据分析系统源码适合具备一定Spark、Kafka基础、希望动手实践流式数据处理的中高级开发者。项目以MQTT协议采集智能家居设备传感器数据经Kafka消息队列实现实时传输再由Spark完成处理分析并写入PostgreSQL同时提供Web实时仪表板与NiFi数据可视化界面覆盖从数据采集、存储到分析展示的完整链路。资源包共16个文件包含zip依赖库、db数据库文件、yml容器编排配置、sql建表脚本、ino设备端程序、py处理脚本及sh启动脚本等压缩包约174KB结构紧凑便于快速部署。目前已有53人学习下载。读者可借此掌握Docker环境搭建、Kafka与Spark协同、HDFS存储及仪表板展示的完整实现思路是理解智能家居实时分析架构的实用参考。1. 智能家居数据管道为什么 Spark 加 Kafka 是绕不开的组合家里几十个传感器门窗磁、温湿度、人体红外、智能插座每秒都在吐数据。单机脚本跑得动一百台设备跑不动一万台批处理能算昨天的用电报表算不了“现在客厅有人且温度高于 28 度就联动开空调”。这就是智能家居数据分析最真实的痛点数据源持续不断、事件有先后顺序、既要实时告警又要离线复盘。Kafka 负责把散落在各处的设备事件可靠地收上来Spark 负责把流式和批量两条线用同一套代码逻辑算清楚。这套组合不是赶时髦而是因为设备消息天然是追加日志Kafka 的分区模型刚好匹配设备 ID 的并行度Spark Structured Streaming 又能让流处理和批处理共享同一份 DataFrame API。适合谁适合手里已经有智能家居设备数据、想从“能存能看”走到“能算能控”的工程师也适合想拿一个完整项目练手 Spark 和 Kafka 协同的数据开发。下面按“先跑通最小链路再谈参数和坑”的顺序拆开讲。2. 从设备消息到 Kafka Topic数据接入层怎么设计2.1 智能家居事件模型与 Topic 划分智能家居的数据不像电商订单那样规整。一个温湿度传感器上报的是{deviceId:th-001,type:temperature,value:26.5,ts:1710000000000}一个人体红外上报的是{deviceId:pir-007,type:motion,value:1,ts:1710000000123}。如果所有设备塞进一个 Topic下游解析时字段对不齐分区也没法按设备隔离。常见做法是按数据大类拆 Topicsmarthome.sensor.telemetry放周期性遥测smarthome.device.event放开关、告警、联动触发这类离散事件。分区数怎么定一个分区只能被一个消费者线程读分区太少限制并行度太多增加协调开销。我一般按“峰值每秒消息数 ÷ 单分区每秒可处理消息数”估算再留一倍余量。比如峰值 5000 条/秒单分区实测能扛 2000 条/秒那 3 个分区够用但为了后续加消费者设 6 个更稳妥。分区键用deviceId这样同一设备的消息严格有序做状态判断时不会因为乱序误判。2.2 用 Python 模拟设备上报并写入 Kafka没有真实设备时先用脚本造数据把链路跑通。下面这段代码模拟 10 个设备每 0.5 秒发一条温湿度或人体感应消息。from kafka import KafkaProducer import json, time, random # 连接本地 Kafkabootstrap_servers 按实际地址改 producer KafkaProducer( bootstrap_servers[localhost:9092], value_serializerlambda v: json.dumps(v).encode(utf-8), # 按 deviceId 做分区保证同一设备消息有序 key_serializerlambda k: k.encode(utf-8) ) devices [fth-{i:03d} for i in range(5)] [fpir-{i:03d} for i in range(5)] while True: for dev in devices: if dev.startswith(th): msg { deviceId: dev, type: temperature, value: round(random.uniform(18, 32), 1), ts: int(time.time() * 1000) } topic smarthome.sensor.telemetry else: msg { deviceId: dev, type: motion, value: random.choice([0, 1]), ts: int(time.time() * 1000) } topic smarthome.device.event # key 传 deviceIdKafka 按 key hash 落到固定分区 producer.send(topic, keydev, valuemsg) producer.flush() time.sleep(0.5)逻辑说明value_serializer把字典转成 JSON 字节流key_serializer把设备 ID 转成字节流。Kafka 对相同 key 的消息会分配到同一分区这是后续做设备级状态计算的前提。producer.flush()在循环里调用是为了确保消息真正发出生产环境可以靠linger_ms攒批但调试阶段直接 flush 更直观。参数上bootstrap_servers如果连的是集群写两三个 broker 地址用逗号隔开不要只写一个。acks默认是 1对智能家居遥测数据够用如果是门锁开关这类事件建议设成all避免 broker 落盘前挂掉丢消息。2.3 验证消息是否落进正确分区写完数据别急着写 Spark先用命令行确认 Topic 和分区分布。# 列出所有 topic kafka-topics.sh --bootstrap-server localhost:9092 --list # 查看 telemetry topic 的分区详情 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic smarthome.sensor.telemetry # 从开头消费 10 条看看内容 kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic smarthome.sensor.telemetry --from-beginning --max-messages 10--describe输出里会显示每个分区的 Leader 和 Replica 分布如果某个分区 Leader 是 -1说明副本没同步上先查 broker 日志。--from-beginning只对还没被消费组提交偏移量的情况有效如果之前有消费者提交过得加--group指定新组名才能从头读。这一步看着简单但很多“Spark 读不到数据”的问题根因是 Topic 名拼错或者消息压根没写进去。3. Spark Structured Streaming 消费 Kafka最小可跑通的流处理3.1 流式读取 Kafka 的 DataFrame 写法Spark 3.x 之后Structured Streaming 读 Kafka 就是spark.readStream.format(kafka)一行的事但参数配不对就会卡在启动阶段。下面是最小可跑通的 PySpark 代码。from pyspark.sql import SparkSession from pyspark.sql.functions import from_json, col from pyspark.sql.types import StructType, StringType, DoubleType, LongType spark SparkSession.builder \ .appName(SmartHomeStreaming) \ .config(spark.sql.shuffle.partitions, 6) \ .getOrCreate() # 定义 JSON 消息的 schema字段要和生产者发的对齐 schema StructType() \ .add(deviceId, StringType()) \ .add(type, StringType()) \ .add(value, DoubleType()) \ .add(ts, LongType()) raw_df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(subscribe, smarthome.sensor.telemetry) \ .option(startingOffsets, latest) \ .option(failOnDataLoss, false) \ .load() # Kafka 的 value 是二进制先转字符串再解析 JSON parsed_df raw_df.select( from_json(col(value).cast(string), schema).alias(data) ).select(data.*) # 简单过滤只看温度高于 28 度的记录 hot_df parsed_df.filter(col(type) temperature).filter(col(value) 28) query hot_df.writeStream \ .outputMode(append) \ .format(console) \ .option(truncate, false) \ .trigger(processingTime10 seconds) \ .start() query.awaitTermination()逻辑说明startingOffsets设成latest表示只读启动之后新到的消息调试时如果想让历史数据也进来改成earliest。failOnDataLoss设false是为了在 Kafka 清理了旧 offset 时不至于让整个流挂掉生产环境要配合监控告警。from_json的 schema 必须和生产者发的 JSON 字段类型严格一致value在生产者那边是浮点数这里用DoubleType如果写成StringType会解析出 null。trigger设 10 秒是为了在控制台看到攒批效果实际部署可以改成processingTime1 minute降低小文件压力。3.2 检查点目录与输出模式的选择Structured Streaming 靠 checkpoint 记录消费进度和状态不配 checkpoint 目录流任务重启后会从头再来或者直接报错。query hot_df.writeStream \ .outputMode(append) \ .format(parquet) \ .option(path, /data/smarthome/hot_temperature) \ .option(checkpointLocation, /data/smarthome/checkpoint/hot_temperature) \ .trigger(processingTime1 minute) \ .start()checkpointLocation必须是一个持久化存储路径本地调试用文件系统可以集群上要换成 HDFS 或对象存储。outputMode三种选择append只输出新行适合落 parquetupdate输出被更新的行适合做聚合后写数据库complete每次输出全量结果只适合结果集很小的聚合。智能家居场景里原始遥测落盘用append房间平均温度这种聚合用update写 Redis 或 MySQL。注意 checkpoint 目录一旦启用不要手动删里面的文件否则流任务恢复时会报 offset 不连续。3.3 用 Spark SQL 做窗口聚合每 5 分钟房间平均温度流处理不只是过滤还要做时间窗口聚合。下面按设备分组算 5 分钟滚动窗口的平均温度。from pyspark.sql.functions import window, avg windowed_df parsed_df \ .filter(col(type) temperature) \ .groupBy(window(col(ts).cast(timestamp), 5 minutes), col(deviceId)) \ .agg(avg(value).alias(avg_temp)) query windowed_df.writeStream \ .outputMode(update) \ .format(console) \ .option(truncate, false) \ .option(checkpointLocation, /data/smarthome/checkpoint/avg_temp) \ .trigger(processingTime1 minute) \ .start()window函数的第一个参数是时间列必须是 timestamp 类型Kafka 消息里的ts是毫秒时间戳用cast(timestamp)转换。窗口长度 5 分钟没有设滑动步长默认就是滚动窗口。outputMode(update)让每个微批只输出有变化的设备聚合结果避免重复写全量。这里有个容易翻车的点如果设备上报时间戳用的是设备本地时间且没同步 NTP窗口会错乱建议在接入层统一用服务端接收时间做事件时间或者至少校验设备时间偏差。4. 离线批处理与流批一体同一套逻辑跑两条线4.1 用 Spark SQL 做历史数据复盘流处理解决实时告警离线批处理解决报表和模型训练。智能家居的离线分析常见需求是按天统计每个房间的用电量、找出温度异常时段、分析人体感应与灯光开启的相关性。这些用 Spark SQL 直接查 parquet 文件就行。# 读昨天落盘的 parquet 数据 batch_df spark.read.parquet(/data/smarthome/telemetry/dt2024-03-01) batch_df.createOrReplaceTempView(telemetry) # 按设备统计当天温度最大值、最小值和平均值 spark.sql( SELECT deviceId, MAX(value) AS max_temp, MIN(value) AS min_temp, ROUND(AVG(value), 2) AS avg_temp, COUNT(*) AS sample_count FROM telemetry WHERE type temperature GROUP BY deviceId ORDER BY avg_temp DESC ).show()逻辑说明parquet 按天分区存储时路径里带dt2024-03-01这种分区目录Spark 读取时会自动做分区裁剪只扫对应目录。createOrReplaceTempView注册临时视图后就能用纯 SQL 表达适合数据分析师直接上手。sample_count用来判断数据完整性如果某个设备当天只有几条记录说明设备离线或网络有问题这种异常要在报表里标出来不能直接拿平均值下结论。4.2 流批一体用同一份 schema 和转换逻辑流处理和批处理最大的重复劳动是 schema 定义和字段清洗。把这两部分抽成独立模块流和批都引用同一份代码。# schema_def.py from pyspark.sql.types import StructType, StringType, DoubleType, LongType TELEMETRY_SCHEMA StructType() \ .add(deviceId, StringType()) \ .add(type, StringType()) \ .add(value, DoubleType()) \ .add(ts, LongType()) def clean_telemetry(df): 统一清洗逻辑过滤空设备 ID转换时间戳 from pyspark.sql.functions import col, to_timestamp return df.filter(col(deviceId).isNotNull()) \ .withColumn(event_time, to_timestamp(col(ts) / 1000))流处理里from_json(col(value).cast(string), TELEMETRY_SCHEMA)批处理里spark.read.schema(TELEMETRY_SCHEMA).json(...)清洗函数两边都调clean_telemetry。这样改字段类型或加过滤条件时只改一处不会出现流和批算出来的结果对不上。我见过一个项目流处理里把温度单位从摄氏度改成了华氏度批处理脚本没同步改导致实时告警和离线报表差了 32 度排查了一下午才发现是两套代码。血泪经验就是能抽公共模块就别复制粘贴。4.3 把聚合结果写回 Kafka 供下游消费算完的结果不一定只落库也可以写回 Kafka 让告警服务、App 推送服务去消费。result_df windowed_df.selectExpr( CAST(deviceId AS STRING) AS key, to_json(struct(*)) AS value ) query result_df.writeStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(topic, smarthome.agg.room_temp) \ .option(checkpointLocation, /data/smarthome/checkpoint/agg_to_kafka) \ .outputMode(update) \ .start()写回 Kafka 时DataFrame 必须包含key和value两列且都是字符串或二进制。to_json(struct(*))把整行转成 JSON 字符串。outputMode用update是因为聚合结果会变下游按 key 做幂等更新。注意写回 Kafka 的 Topic 分区数最好和上游一致否则可能出现下游消费并行度不匹配。另外如果下游是告警服务消息里要带窗口起止时间不然收到“平均温度 30 度”却不知道是哪 5 分钟的数据。5. 避坑与排查智能家居流处理里最容易翻车的 4 个点5.1 现象Spark 流任务启动后一直没输出控制台空白原因startingOffsets设成了latest而生产者还没开始发数据或者 Topic 名写错导致订阅了一个空 Topic。另一个常见原因是trigger设了较长的 processingTime比如 5 分钟启动后要等第一个微批结束才输出。解决调试阶段把startingOffsets改成earliesttrigger改成processingTime10 seconds同时用kafka-console-consumer.sh确认 Topic 里确实有消息。如果 Topic 有消息但 Spark 读不到检查subscribe参数有没有拼写错误Kafka 的 Topic 名是大小写敏感的。5.2 现象JSON 解析后字段全是 null原因from_json的 schema 和实际 JSON 字段类型不匹配。比如生产者发的value是整数26schema 里定义成DoubleType通常能兼容但如果生产者发的是字符串26.5schema 里定义成DoubleType就会解析失败返回 null。解决先用kafka-console-consumer.sh看原始消息长什么样再对照 schema 逐字段核对。一个省事的办法是先用StringType把所有字段读进来再用cast转换这样解析失败时能看到原始字符串方便定位。5.3 现象流任务跑一段时间后报 OffsetOutOfRange原因Kafka 的log.retention.hours默认是 168 小时如果流任务停了超过这个时间之前提交的 offset 对应的消息已经被清理重启时就会报 offset 越界。解决把failOnDataLoss设成false让任务跳过丢失的数据继续跑同时检查 checkpoint 目录是否被误删。长期方案是调大 retention 时间或者让流任务保持运行不要长时间停机。如果业务允许丢数据startingOffsets设latest重新开始也行但要评估丢失窗口对业务的影响。5.4 现象窗口聚合结果延迟很高告警不及时原因trigger的 processingTime 设得太长比如 10 分钟加上窗口长度 5 分钟最坏情况下要等 15 分钟才输出。另一个原因是spark.sql.shuffle.partitions设得太大小数据量下反而增加调度开销。解决实时告警场景把trigger降到 10 到 30 秒窗口长度按业务容忍度调整。shuffle.partitions在本地调试时设成 6 到 12 就够集群上按数据量调一般设成 CPU 核数的 2 到 3 倍。如果还是慢看 Spark UI 里哪个 stage 耗时最长通常是 Kafka 读取或 JSON 解析阶段可以增加maxOffsetsPerTrigger限制每批拉取量避免单批数据过大。6. 进阶技巧用水位线处理迟到数据和验证结果一致性智能家居设备网络不稳定消息迟到是常态。Structured Streaming 的水位线机制可以容忍一定程度的迟到同时控制状态大小。下面在窗口聚合基础上加 10 分钟水位线。from pyspark.sql.functions import window, avg, col windowed_df parsed_df \ .withWatermark(event_time, 10 minutes) \ .filter(col(type) temperature) \ .groupBy(window(col(event_time), 5 minutes), col(deviceId)) \ .agg(avg(value).alias(avg_temp))withWatermark必须用在事件时间列上且要在聚合之前调用。10 分钟水位线意味着 Spark 会等 10 分钟之后到达的、事件时间早于水位线的数据会被丢弃。水位线设太长状态存储压力大设太短迟到数据丢得多。我一般先按设备网络质量估算最大迟到时间再留一倍余量。比如 4G 设备最坏迟到 3 分钟水位线设 6 到 10 分钟。验证结果一致性有个笨但有效的办法把同一时间段的数据分别用流处理和批处理跑一遍对比聚合结果。流处理输出到 parquet 后用批处理读同一份原始数据算一遍两边按deviceId和窗口起止时间 join看avg_temp差值是否在浮点误差范围内。如果差得多先查水位线是不是把迟到数据丢了再查流处理的 checkpoint 是否从正确 offset 开始。这个对账步骤在项目上线前跑一次能提前发现大部分逻辑错误。最后说个习惯我每次改完流处理代码不会直接上生产而是先用earliest从头消费一小段历史数据确认输出符合预期后再切latest部署。这个习惯帮我省掉了至少三次半夜爬起来回滚的麻烦。希望帮到你。本文还有配套的精品资源点击获取
返回列表