news 2026/10/7 6:29:13

Spark Kafka智能家居数据分析实战:实时流处理架构与工程实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Spark Kafka智能家居数据分析实战:实时流处理架构与工程实现

简介:一套面向智能家居场景的完整数据分析源码包,以Spark和Kafka为核心,配合MQTT、Docker等技术,适合有一定大数据基础、正在做课程设计或毕业设计的开发者。项目打通从设备数据采集到实时仪表盘展示的完整链路,包含MQTT协议接入、Kafka消息队列转发、HDFS分布式存储、Spark处理分析并落库PostgreSQL、NiFi可视化等关键环节,对理解物联网数据管道整体设计很有参考价值。资源包共16个文件、约174KB,主要包含Arduino设备端源码及DHT传感器库、Mosquitto与Kafka相关配置、Spark处理脚本与提交脚本、PostgreSQL建表SQL、Docker编排文件、Hadoop环境变量配置、NiFi数据库文件,以及README安装使用说明,目录结构清晰,便于按模块查阅和复用。当前页面已有53人学习浏览,适合作为物联网大数据项目的入门参考和快速起步模板,尤其适合用作选题方案、项目演示和二次开发底稿。

1. 基于Spark和Kafka的智能家居数据分析系统:这套源码到底在解决什么

不少人在做智能家居数据分析这个方向时,最尴尬的局面是:设备数据已经源源不断产生,却不知道该往哪里放、怎么算。直接写 MySQL,热点家庭的数据一上来就把连接池打满;不写库,几十个传感器上报的温湿度、能耗、烟感状态又只能在日志里躺着。这套基于 Spark 和 Kafka 的智能家居数据分析系统,核心就是把 Kafka 当成削峰填谷的缓冲层,用 Spark 做实时清洗和窗口聚合,最终把设备原始数据变成能耗报表、异常告警这类可用的分析结果。

它适合三类人:刚学完 Spark 和 Kafka 基础、想找一个能跑通全链路的实战案例的;要交课程设计或毕业设计、需要一个完整工程的;以及公司内部想低成本搭一套设备数据实时分析平台、又不打算从零写架构的。整套系统用到的组件不多,但链路真实,照着源码跑通之后,换数据源、换存储、加告警规则都是顺手的事。

2. 架构先立住:为什么是 Kafka 削峰、Spark 计算

2.1 数据链路怎么走:设备端到 Spark 再到存储

先看数据特征。智能家居场景里,设备类型和上报频率差异极大:门磁一天可能只有几条事件,能耗表几秒一条,烟感传感器平时静默、告警时突然井喷。如果让这些数据直接打向数据库,写放大非常严重,峰值时段还会拖垮存储。Kafka 在这条链路里起的核心作用是削峰填谷,所有设备数据先进同一个 topic,下游 Spark 按自己的节奏消费,无论上报多么汹涌,存储侧都不会被瞬时流量冲垮。

链路本身不复杂,常见做法是:模拟器或真实网关把数据按 JSON 格式发送到 Kafka 的 dps-event 主题,Spark Structured Streaming 从该主题消费,经过清洗、窗口聚合后,把明细和统计结果分别写入 MySQL、Redis 或 Parquet。拿到源码包后,先按模块把工程拆开看,通常代码结构是下面这张表的样子:

模块职责核心入口主要配置
common公共配置与工具类ConfigLoaderbroker 地址、topic 名、连接池
producer设备数据模拟器DeviceDataProducer上报频率、设备数量、JSON schema
spark-streaming流处理主程序SmartHomeStreamingApp窗口时长、水位线、输出模式
storage结果落库与查询写 MySQL/Redis 的 DAO表结构、连接串、批量参数

为什么非要 Kafka 加 Spark,而不是 MQTT 直连 Spark 或者数据直接写 Redis?MQTT 更适合控制指令的下发,消息一旦被消费就没了,无法回放;而直接把原始值塞进 Redis,等于把分析逻辑全堆在业务代码里,Redis 内存也扛不住长期明细。Kafka 的 offset 机制给了数据一次“后悔药”——消费程序挂了、逻辑写错了,只要不删 topic,随时可以从最早偏移重新消费。Spark 这边,Structured Streaming 对 Kafka Source 的支持最成熟,offset 由 Spark 内部管理,配合 checkpoint 能做到断点续跑,这套组合的可靠程度是经过大规模生产验证的。

2.2 选型理由:Kafka 吞吐量和 Spark 流处理的边界在哪

为什么用这对组合而不是 Flink,是很多人在网上搜 spark 数据分析案例时最容易纠结的问题。Flink 的吞吐和延迟指标确实更漂亮,但学习成本和部署运维成本也更高。智能家居分析场景里,很少需要毫秒级响应,Spark 的微批模型完全够用。更重要的是,Spark 能顺带做批处理:同一个聚合逻辑,实时跑就挂在 Structured Streaming 上,复盘昨天的数据就换成 Spark SQL 直接读 Parquet,批流一体省掉一套代码。如果项目里同时有日报、月报这类离线需求,选 Spark 的收益会明显大于 Flink。

Kafka 侧最容易翻车的是 partition 设计。常见误用是给每个设备建一个 topic,分区数几百上千,最后 Kafka 集群文件句柄爆炸,Spark 消费也跟不上。智能家居场景不需要按设备粒度开 topic,统一一个 dps-event,按 home_id 或 device_type 做 key 就够了。这样同一个家庭的数据落到同一分区,下游按家庭做窗口聚合时,数据乱序程度会轻很多。另一个必须核对的是版本匹配:Spark 3.x 通常搭配 Kafka 2.8 以上和 Scala 2.12/2.13;如果源码 pom 里写的是 Spark 2.4,就别直接连 Kafka 3.x 的 broker,老版本 kafka-clients 的协议兼容性会让你在连接阶段就卡住。

2.3 拿到源码 zip 后,从哪几个文件开始读

不要急着点运行。先把目录结构摊开,找到 pom.xml 或 build.sbt,核对三个版本号:spark-core 的版本、scala.binary.version、kafka-clients 版本。这三个数对齐了,再往下走。接着看配置文件,一般集中在 application.conf 或 kafka.properties 里,里面有 broker 地址、topic 名、消费组 ID、checkpoint 路径。把这些值先对照自己的环境改一遍,再动代码。

读代码的顺序建议是:先读 producer 主类,搞清楚数据长什么样;再读 spark-streaming 主类,看它订阅哪个 topic、解析哪些字段、输出到哪;最后读 storage 相关代码,确认结果表的 schema。Spark 消费端代码量通常最多,抓住三个点就够:readStream 的 format 和 options、from_json 的数据结构定义、输出 sink 是写 MySQL 还是 Redis。读通这三处,整个系统就从黑匣子变成可以随意改的工程。

本地搭建时,不需要一开始就搭完整 Spark 集群。装好 JDK 8+、Kafka 2.8+、Spark 3.x,Kafka 如果不想配 ZooKeeper 就用 KRaft 单机模式,Spark 用 local 模式,先把链路通了再说。有一点要注意:如果源码是在 Linux 环境写的、你却在 Windows 上跑,需要提前配好 winutils 并设置 HADOOP_HOME,否则 Spark 启动会直接报找不到 winutils.exe,这是 Windows 本地跑 Spark 的经典坑。

3. 把设备数据喂进 Kafka:模拟数据生成与 Producer 参数

3.1 先定数据格式:字段、单位、时间戳怎么设计

模拟数据是整个系统能不能跑出真实感的基础。常见的做法是定义这样一个 JSON 结构:

{ "home_id": "H10023", "device_id": "D30048", "device_type": "temperature", "value": 26.8, "unit": "celsius", "event_time": "2024-11-20T21:15:03+08:00" }

字段设计上有几个细节值得说。home_id 是聚合的第一维度,所有按家庭的统计都靠它;device_id 要全局唯一,用来关联设备和房间位置;device_type 是业务维度,聚合和告警规则都以它分类;value 和 unit 分开存放,方便在 Spark 清洗阶段做单位统一和异常值过滤;event_time 一定要用带时区偏移的 ISO8601 字符串,而不是裸的时间戳数字——分布式环境里,裸时间戳在时区转换上非常容易出错,而且出了问题极难排查。

模拟器生成数据时还要维护一份设备类型枚举,常见的有 temperature、humidity、smoke、door_magnet、energy_meter。每种类型的 value 范围不一样:温度一般 15 到 35 摄氏度,湿度 40 到 70,能耗表 0 到 5000 瓦。模拟器按类型给范围和随机分布,下游聚合出来的结果才像真的。如果把温度写成 30 万、能耗写成负数,后续所有清洗逻辑都会失真。

3.2 写 Producer:acks、linger.ms、batch.size 三个参数决定吞吐

模拟数据的 Producer 用 Java 写最通用,核心代码如下:

import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.ProducerConfig; import java.util.Properties; public class DeviceDataProducer { public static void main(String[] args) throws InterruptedException { Properties props = new Properties(); // 单机跑就是 localhost:9092,集群环境填多台 broker 地址 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); // 三个核心吞吐参数:acks、linger.ms、batch.size props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.LINGER_MS_CONFIG, 100); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); // 如果单条数据超过 1M(比如带了监控图),要同步调大客户端和 broker 限制 props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 10485760); KafkaProducer<String, String> producer = new KafkaProducer<>(props); // 模拟 10 个家庭、每个家庭 20 个设备,共 200 个设备轮流上报 for (int i = 0; i < 100000; i++) { String homeId = String.format("H%05d", i % 10); String deviceId = String.format("D%05d", i % 200); String json = "{\"home_id\":\"" + homeId + "\",\"device_id\":\"" + deviceId + "\",\"device_type\":\"temperature\",\"value\":25.5,\"unit\":\"celsius\"}"; // key 传 homeId,保证同一家庭的数据进同一分区 producer.send(new ProducerRecord<>("dps-event", homeId, json)); Thread.sleep(50); } producer.close(); } }

逻辑说明:key 传 homeId 是为了让同一家庭的数据落到同一分区,下游按家庭做窗口聚合时数据乱序程度会小很多。循环里的 sleep 50ms 是在模拟真实设备的上报节奏,而不是满速打满 broker,这样跑出来的消费延迟才贴近现实。

参数说明:acks=all 表示等所有副本确认,配合默认的 retries=3,能有效避免“发出去但丢了”的尴尬。linger.ms=100 是让 Producer 攒 100 毫秒的批再发,如果设成 0,每条消息立刻发出,吞吐反而最差;batch.size 是单批字节上限,这两个参数必须一起调,只调一个看不出效果。max.request.size 的默认值是 1MB,很多搜“kafka 接收 1m”的人遇到的就是这个限制——如果上报的数据里带了图片或音频片段,必须调大这个值,同时 broker 侧要同步调 message.max.bytes 和 replica.fetch.max.bytes,否则客户端不报错、broker 直接拒收。

还有一个容易被忽略的点:把 linger.ms 调到 500 甚至 1000 确实能提升吞吐,但会带来可见的消息延迟。温度、能耗这类指标对延迟不敏感,100 毫秒是稳妥值;但烟感和紧急按钮的消息如果也走同一个 topic,最好单独建一个 alert-topic,把 linger.ms 设为 0,保证告警消息不加缓冲地立即发出去。这是我在真实项目里常用的做法。

3.3 数据有没有进 Kafka:控制台消费与可视化工具确认

数据发完之后,先别急着写 Spark,确认数据进了 topic 再说。用 Kafka 自带的命令行脚本验证最快:

# 先看 topic 有没有建立 kafka-topics.sh --bootstrap-server localhost:9092 --list # 没有就建一个,本地测试 1 个分区够用,正式跑建议 3 个分区 kafka-topics.sh --bootstrap-server localhost:9092 \ --create --topic dps-event --partitions 3 --replication-factor 1 # 从最早偏移读 10 条,确认 JSON 格式和字段名 kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic dps-event --from-beginning --max-messages 10

先 list 确认 topic 存在,不存在就建。分区数本地设为 1 就行,正式环境至少 3 个分区,这样 Spark 的并行度才能拉起来。--from-beginning 表示从最早偏移开始读,加 --max-messages 10 只读 10 条,避免刷屏。

命令行验证没问题后,可以顺手用一个 Kafka 可视化工具(常见的有 kafka-ui、kafka-map 这类开源项目)看一下 topic 的分区 leader、lag、消费组进度。这些工具对快速定位“数据卡在 Producer 还是消费端”特别有用:lag 持续上涨说明生产快消费慢,问题大概率在 Spark 端处理能力不足;lag 为 0 但迟迟不出结果,问题就在窗口触发或写库逻辑。这个判断能省掉大量排查时间。

4. Spark 侧的核心计算:清洗、统计与实时告警

4.1 Structured Streaming 还是老版 Spark Streaming:边界判断

源码里如果给的是老版 Spark Streaming(DStream),我的建议是直接重写成 Structured Streaming。两者的差异用一张表看更清楚:

对比项Spark Streaming (DStream)Structured Streaming
编程模型RDD 流,手动管理状态无界表,像批处理一样写 SQL
时间语义处理时间为主原生支持 event time 和 watermark
offset 管理手动维护或依赖 checkpointSpark 自动管理,配合 checkpoint 更可靠
故障恢复可能重复消费配合幂等 sink 实现恰好一次
适用场景老项目维护新项目首选

Structured Streaming 把流当成一张不断追加的无界表,写窗口聚合、join 维表时和批处理 SQL 几乎一样,代码量少一半。真正生产环境里的“恰好一次”语义,也更容易通过 readStream 加 foreachBatch 加幂等写实现。如果这套智能家居系统的源码用的是 DStream,重写成本并不高,换掉读取和输出两层就够了,业务逻辑几乎原样平移。

4.2 从 Kafka 读数据并做 JSON 解析与清洗

很多人在搜“spark 中读取 json”时没意识到,流式读 JSON 的写法跟批处理不一样。批处理可以直接 spark.read.json,流式必须先读 Kafka 的 value 字段,再手动做 from_json 解析。核心代码:

// Scala 版,Spark 3.x + Kafka 依赖 import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ val schema = new StructType() .add("home_id", StringType) .add("device_id", StringType) .add("device_type", StringType) .add("value", DoubleType) .add("unit", StringType) .add("event_time", StringType) val raw = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "dps-event") .option("startingOffsets", "earliest") .option("maxOffsetsPerTrigger", "10000") .load() val deviceEvents = raw .selectExpr("CAST(value AS STRING) as json_str") .select(from_json(col("json_str"), schema).as("data")) .select("data.*") .filter(col("value").isNotNull && col("value") > 0) .withColumn("event_time_ts", to_timestamp(col("event_time"), "yyyy-MM-dd'T'HH:mm:ssXXX"))

逻辑说明:from_json 给数据定义显式 schema,避免 Spark 推断类型时把数值列当成字符串;filter 在流内就把 null 值和非法值剔掉,不要等写入下游再清理,那样既浪费存储又污染报表;event_time_ts 把带时区偏移的字符串转成 timestamp,后续窗口聚合和 watermark 都依赖这个字段。

参数说明:startingOffsets=earliest 表示从最早偏移开始读,适合模拟数据重跑场景;maxOffsetsPerTrigger 默认是尽量多读,本地小机器建议设成 10000,防止一次拉太多把 driver 内存打爆。from_json 遇到解析失败的行会默认丢弃,如果不想粗暴丢,可以在 select 之后检查 corrupt 字段,单独统计脏数据率,这样能保留“哪些设备总发坏数据”的线索,比直接 dropMalformed 更实用。

如果消费端经常出现“kafka 消息延迟高”,且所有逻辑都正常,大概率是 maxOffsetsPerTrigger 没设边界或分区数太少。把分区数提到与 Spark executor 核数匹配,加上合理的 maxOffsetsPerTrigger,lag 自然降下来。另外,Structured Streaming 的消费组 ID 由 checkpointLocation 管理,旧版代码里手动指定 group.id 的方式常常造成重复消费,这是迁移时最容易忽略的差异。

4.3 5 分钟窗口聚合:能耗统计和烟感告警怎么做

窗口聚合是整套系统的业务核心。能耗统计按 5 分钟滚动窗口、按家庭聚合;烟感告警用 1 分钟窗口统计报警次数,连续触发才告警:

// 5 分钟滚动窗口 + 按家庭维度聚合能耗 val energyStats = deviceEvents .filter(col("device_type") === "energy_meter") .withWatermark("event_time_ts", "10 minutes") .groupBy( col("home_id"), window(col("event_time_ts"), "5 minutes") ) .agg(sum("value").as("total_energy_kwh"), avg("value").as("avg_power_w")) // 烟感异常:value=1 表示报警,1 分钟内连续 3 次才触发 val smokeAlert = deviceEvents .filter(col("device_type") === "smoke" && col("value") === 1) .withWatermark("event_time_ts", "5 minutes") .groupBy(col("home_id"), window(col("event_time_ts"), "1 minutes")) .count() .filter(col("count") >= 3)

window 函数会生成 window_start 和 window_end 两列,groupBy 里带上 home_id 就完成了按家庭和时间的双重维度聚合。withWatermark 的参数和窗口长度强相关:超过水位的迟到数据会被丢弃,水位一般设为窗口长度的两倍,5 分钟窗口给 10 分钟水位。烟感告警加“1 分钟连续 3 次”这个条件,是为了过滤瞬时误报——很多烟雾报警器会因为蒸汽短促触发一次,如果每次都告警,运维很快就会疲劳。

输出模式怎么选也值得说。写 MySQL 这类目标表时用 append 模式配合 foreachBatch,因为 append 只在窗口结束后输出一批数据,批次语义清晰;做实时大屏用 update 模式,它只推送更新的行,减小下游写压力;complete 模式会把所有聚合结果全量输出,家庭数少的时候没问题,到几十万家庭时直接把 driver 拉垮。

4.4 结果写库:MySQL、Redis、ES 如何分工

原始明细、统计结果、实时告警这三类数据不适合放进同一个存储。我一般这么分:清洗后的明细写 Parquet 离线目录,供后续 Spark SQL 复盘和日报;5 分钟聚合结果写 MySQL,供业务查询;大屏实时值写 Redis,用 home_id 做 key。Structured Streaming 允许同一个 readStream 派生多个输出流,一个写 Redis、一个写 MySQL,互不干扰。

MySQL 写入的标配写法是 foreachBatch:

import org.apache.spark.sql.{DataFrame, SaveMode} energyStats.writeStream .foreachBatch { (batchDF: DataFrame, batchId: Long) => batchDF.withColumn("batch_id", lit(batchId)) .write .mode(SaveMode.Append) .jdbc("jdbc:mysql://localhost:3306/smart_home", "energy_stats_5min", props) } .option("checkpointLocation", "/tmp/chk/energy_stats") .outputMode("append") .start()

逻辑说明:foreachBatch 把每个微批当成一个 DataFrame,用 JDBC 批量写入,比逐条 insert 快两个数量级。batch_id 列用来对账,下游发现某个窗口重复计算时,可以按 batch_id 剔除。

参数说明:checkpointLocation 必须一对一绑定,同一个路径不能同时给两个流用,否则 offset 冲突。SaveMode.Append 配合目标表的联合唯一键(home_id + window_start)实现幂等,重复跑批也不至于让统计数字翻倍。写入 MySQL 的 JDBC driver 要提前放到 Spark 的 classpath 里,本地跑常常报 ClassNotFound,就是这个依赖缺失导致的。

5. 避坑清单:从编译到结果落库的 5 个真实事故

5.1 现象:Spark 一直消费不到新数据,控制台却能看到消息在发

原因:Structured Streaming 默认从 topic 最新偏移开始读,而每次重启都会创建新的消费组。如果先启动 Producer 灌数据、再启动 Spark,第一批数据早就发完了,Spark 从 latest 开始等待新消息,于是看起来“一直没反应”。源码里如果把 startingOffsets 写死了 latest,这个现象几乎必现。

解决:本地测试把 startingOffsets 改成 earliest,先启动 Spark 再启动 Producer;线上环境不要每次重启都从头消费,让 checkpoint 管理偏移,正常情况下不需要手工改这个参数。

5.2 现象:5 分钟窗口算出来的能耗数值和手工对账差一大截

原因:设备端事件时间用了本地时区,Spark session 默认 timeZone 是 UTC,window 函数按 UTC 切窗口,数据全被切到了相邻窗口。另一个常见原因是时间字符串格式没对上,to_timestamp 解析失败,整行被静默丢弃,统计自然偏少。

解决:统一约定——模拟器时间用 ISO8601 带偏移字符串,Spark 侧设置 spark.conf.set("spark.sql.session.timeZone", "Asia/Shanghai"),窗口按设备本地时间切。解析失败的行不要静默丢,单独统计脏数据比例,观察一段时间后就能定位是哪个设备型号的格式不兼容。

5.3 现象:Kafka 连接超时,localhost 跑得好好的,换服务器 IP 就挂

原因:Kafka broker 的 advertised.listeners 没配置。broker 启动时把自己的 advertised 地址写进元数据,客户端按这个地址去连;默认值正好是 localhost,所以本地没问题,换到服务器后客户端拿到的还是 localhost,自然连不上。这个坑在 Spark 集群提交作业时高频出现。

解决:broker 配置里显式写 advertised.listeners=PLAINTEXT://服务器实际IP:9092,并确认安全组对这个端口放行。容器环境还要区分容器内地址和宿主机地址,两个都得能通,否则会出现“同机可以连、跨机必超时”的诡异现象。

5.4 现象:作业跑一会儿 OOM,driver 或 executor 内存被打满

原因:两个最可疑的地方。一是 from_json 解析的原始 JSON 里嵌套了数组或大字符串,数据集体积比预期大得多;二是聚合流没设 watermark,或水位设置比窗口还短,state store 里堆积了所有历史窗口,update 模式频繁写 checkpoint,内存和磁盘一起吃紧。

解决:先把并行度调到合理区间,小数据量把 spark.sql.shuffle.partitions 从默认 200 改成 16;给聚合流补上 withWatermark,水位至少是窗口长度的两倍;最后再考虑加 spark.executor.memoryOverhead。OOM 排查不要一上来就加内存,先看状态存储,只有把 state 保留周期压下来才能治本。

5.5 现象:MySQL 写入慢,每秒几百行都勉强度日

原因:foreachBatch 内部用了逐行插入,或者 JDBC 连接串没加批量参数。Spark 侧通过 JDBC 直连 MySQL,逐行插入一次握手一次,延迟全消耗在网络上。源码里如果写的是 collect 然后 for 循环插入,基本就是这个毛病。

解决:把写入逻辑改成 batchDF.write.jdbc,连接串加 rewriteBatchedStatements=true 和 useServerPrepStmts=true,批量写性能能提升一个数量级。同时给目标表建好联合索引 home_id + window_start,否则数据量上来之后查询也会跟着变慢。

6. 让系统更像生产:配置外置、断点续跑与结果对账

代码里写死 broker 地址和 topic 名,换环境就得改代码重新编译,这是新手最常干的事。我一般把连接参数集中到一个 application.conf,按 dev 和 prod 分两套 profile,启动时指定环境即可:

kafka { bootstrap.servers = "localhost:9092" topic = "dps-event" } spark { shuffle.partitions = 16 timezone = "Asia/Shanghai" } checkpoint { dir = "/tmp/chk/smart-home" }

checkpoint 目录是整套系统的后悔药。想重算今天的数据,不要删掉 checkpoint 整个目录再重启,那样会丢失消费进度;正确做法是换一个新的 checkpoint 路径并配合 startingOffsets=earliest,让旧路径保留现场,新路径从头算,两份结果可以做交叉对账。如果想回滚,把旧的 checkpoint 路径配回去,消费偏移自然回到上次记录的位置。

我最早做这套系统时,图省事把窗口时长写死成 30 秒,结果所有聚合结果都差半小时,后来统一改成从配置里读,才彻底消停。跑数据平台的工程问题不是算法难,而是数据链路里任何一个环节静默出错都会让结果失真,所以最后一步永远是“数对不对得上”:从 Kafka 消费端记录总条数,再对 MySQL 里每类设备的聚合数求和,做小时级别的差值校验。把这个校验写成 Spark SQL 离线任务每天早上定时跑,比肉眼盯监控靠谱得多。希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/7 6:29:13

AI编程智能体实战:从工具选型到提示词设计的完整指南

过去半年&#xff0c;我身边几乎所有写代码的朋友都在聊同一件事&#xff1a;AI编程智能体。有人焦虑&#xff0c;觉得初级程序员的活快被干完了&#xff1b;也有人兴奋&#xff0c;说自己已经让AI把一个两周的活压缩成了两天。就我实际测下来的感受&#xff0c;这两种说法都不…

作者头像 李华
网站建设 2026/10/7 6:28:58

果蔬图像分类实战:4200张标注数据集与PyTorch ResNet训练全指南

简介&#xff1a;一份面向图像分类任务学习与算法验证的常见果蔬多类别数据集&#xff0c;覆盖香蕉、苹果、梨、葡萄、橙子、胡萝卜、辣椒、洋葱、土豆等36个类别&#xff0c;共约4200张已标注图像。数据经过统一预处理&#xff0c;可直接作为分类网络输入&#xff0c;并已划分…

作者头像 李华
网站建设 2026/10/7 6:28:28

S7-1200控制步进电机实战:从选型接线到博途组态与调试

去年接了一个传感器装配工装的项目&#xff0c;转盘旁要加一台步进电机驱动的推料机构&#xff1a;正转把物料推到检测位&#xff0c;检测完再反转退回来&#xff0c;转速要在触摸屏上能调。用西门子S7-1200做PLC控制步进电机&#xff0c;听起来无非是梯形图加脉冲输出&#xf…

作者头像 李华
网站建设 2026/10/7 6:27:55

Flink读取Kafka实战:从环境搭建到参数调优与避坑指南

简介&#xff1a;一份基于 Apache Flink 的实时数据处理工程包&#xff0c;完整演示从 Kafka 消费数据、执行流式计算、再将结果写入 Redis 集群与 MySQL 的链路。面向大数据开发初学者及有实时数仓落地需求的工程师&#xff0c;适合用于学习 Flink Connector 配置、算子应用和…

作者头像 李华
网站建设 2026/10/7 6:27:51

JSP+SQL实现交通信息管理系统实战指南

简介&#xff1a;本资源是一套完整的基于JSP与SQL技术的本科毕业设计项目材料&#xff0c;面向计算机专业学生、Web开发初学者及Java后端入门学习者&#xff0c;聚焦智能交通信息管理这一典型B/S架构应用场景。压缩包共含项目报告、源代码、开题报告、答辩PPT与外文翻译等核心文…

作者头像 李华
网站建设 2026/10/7 6:27:21

FPGA实现10G/25G UDP线速网络栈的工程落地路径

1. 这不是“又一个UDP demo”&#xff0c;而是一套能跑满线速的FPGA网络栈落地路径你有没有试过在Vivado里拖一个UDP IP核&#xff0c;写几行Verilog发个包&#xff0c;然后用Wireshark抓到数据——心里一热&#xff0c;以为成了&#xff1f;结果一上真实流量&#xff0c;ping延…

作者头像 李华