ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于Spark与Kafka的智能家居流式数据处理实践

基于Spark与Kafka的智能家居流式数据处理实践 简介这是一套基于Spark和Kafka的智能家居数据分析系统源码面向物联网开发者、大数据初学者以及需要构建实时数据处理管线的工程师可用于课程设计、毕业设计或技术实战。系统采用MQTT协议从智能家居设备采集环境与状态数据以Kafka消息队列作为中转保障实时可靠性由Spark执行处理与分析结果存入PostgreSQL同时通过NiFi进行可视化展示整体覆盖数据采集、传输、存储、计算到呈现的完整流程。压缩包共16个文件大小约174KB包含Arduino主控程序与传感器库、Docker编排配置、数据库建表脚本、PySpark处理脚本、任务提交脚本及说明文档文件类型涉及zip、ino、yml、env、conf、sql、db、py、sh、md等能支撑从设备端固件到服务端部署的关键环节。目前已有53人学习下载适合按README与docker-compose快速搭建运行环境复现实时仪表板和统计结果帮助理解物联网数据流与大数据组件协作方式。1. 基于Spark和Kafka的智能家居数据分析系统它到底在解决什么问题做智能家居数据分析最典型的场景不是“数据量有多大”而是“数据来得太快、太杂、太乱”。一个家庭里温湿度传感器、门窗磁、人体红外、智能插座、烟雾报警器同时上报每秒可能就产生几十到几百条记录高峰期像晚上六点到十点所有设备都在动。传统的做法是把数据先落库再跑分析但到了这个量级入库太慢、查询太慢、报表刷新太慢整个链路处处是瓶颈。这个基于Spark和Kafka的智能家居数据分析系统本质上是把“数据先放消息队列缓冲再用流处理引擎实时消费分析”的经典大数据架构落到了智能家居这个具体业务上Kafka负责扛住写入洪峰Spark负责把数据洗成可分析的形状Redis或MySQL承接最终结果供前端和大屏使用。适合的人群很明确正在做物联网或智能家居数据平台、想把流处理架构跑通但又不想从零造轮子的人以及准备拿真实项目练手Spark和Kafka的工程师。2. 整体架构与数据流转为什么一定是Kafka加Spark这对组合2.1 智能家居数据的写入特征决定了必须先有缓冲层智能家居设备上报数据有一个很反直觉的特征总量不大但瞬时峰值很高。一套两室一厅的房子可能只有三十来个传感器每个传感器每5秒上报一条平均每秒才6条数据单看这个量级根本没有压力。但如果是整栋楼、整个小区呢上千户家庭同时活跃设备种类从温湿度传感器到智能门锁、烟雾报警器、水浸传感器上报频率各不相同峰值流量可能比平均值高出十倍以上。更麻烦的是设备上报行为不可控——某个品牌的中枢网关重启后底下十几个子设备会同时补报历史状态瞬间打出一波与当前时刻无关的旧数据。这种写入特征决定了架构选型的第一条原则不能让分析引擎直接对接设备端。如果让Spark直接接收设备上报那Spark的微批调度反而会被这种无规律的突刺打乱节奏而且一旦Spark任务重启或升级设备端没有重试机制数据就直接丢了。所以中间必须有一层既能扛峰值、又能持久化消息的缓冲Kafka在这个位置几乎是没有争议的选择。Kafka的partition机制天然支持水平扩展吞吐量随机器数量线性增长而且消息持久化在磁盘上消费者挂了之后从offset恢复消息一条都不会少。2.2 为什么实时分析层选Spark而不是直连Kafka消费Kafka本身不提供分析能力它只有生产、存储、消费这三件事。数据到了Kafka之后接下来需要一个能持续消费、做清洗、做聚合、写结果的计算引擎。这里能选的东西其实不少Storm太老Flink和Spark Structured Streaming是当前的两个主流方向。这个项目选了Spark核心原因是技术栈收敛和生态成熟度Spark的Dataset/DataFrame API对SQL用户极其友好一段流处理逻辑可以先用批处理模式调通再切到流处理代码几乎不用改。这一点做物联网项目时太重要了因为智能家居数据里大量是格式不规整的JSON用Spark SQL解析、过滤、Join比用Flink的DataStream API写起来直观得多。另外还有一个很实际的考量流处理跑起来之后十有八九还要对历史数据做离线分析——比如统计某个户型一周的用电规律、分析传感器上报频率与设备异常的关系这些是典型的批处理任务。一套Spark集群同时支撑实时流和离线批运维成本比维护两套引擎低一半。2.3 核心数据模型从设备原始报文到分析宽表在写任何代码之前先把消息格式定死这是这个项目最值得参考的设计决策。Kafka里的每一条消息就是一个JSON字符串包含设备ID、时间戳、数据类型和数据体四部分。我的建议是消息里不嵌套任何业务语义只做最原始的透传。设备端上报是什么样进Kafka就是什么样分析层的清洗逻辑不在生产端做。这样做的好处是事后排查问题时有据可查Kafka里保留的是最原始的数据就算Spark作业洗错了重新从Kafka消费一遍就能修复不必回头找设备端要数据。Kafka的topic按设备类型分区设计一个topic对应一类设备比如sensor_temperature、sensor_motion、device_switch每个topic根据预期峰值流量设6到12个partition。这里要特别注意topic数量一旦定了就不要频繁改改partition数会打乱既有消息的key分布导致部分分区数据倾斜。宁可前期多估一点流量多建几个partition也不要上线后再扩容。3. 环境与版本选型Spark集群和Kafka集群的搭建及版本匹配要点3.1 版本兼容是第一个要命的坑入手这套系统之后第一件事不是写代码而是确认Spark和Kafka的版本兼容关系。这个坑几乎每个做流处理的人都踩过而且踩法五花八门。Spark Structured Streaming通过Spark的kafka connector消费Kafka这个connector的编译版本必须和Kafka broker端兼容。Kafka 3.x早期版本和Spark 3.0以下的组合经常出现反复重试、offset提交失败、消费延迟居高不下这些问题——绝大多数情况下不是业务代码写错了而是客户端协议版本对不上。我一般会直接看两个版本号Spark的scala版本和Kafka的客户端协议版本。比如Spark 3.2以上官方默认支持Kafka 2.8的协议而Kafka 3.0以后的broker对旧协议做了大量兼容处理但反过来Kafka 3.4的某些新特性老版本的Spark connector根本不认识。保守做法是Kafka用3.x系列Spark也选3.x系列两者都在同一年代版本内避免跨越太大。3.2 Kafka集群安装用KRaft模式省掉ZooKeeperKafka的部署方式这两年发生了很大变化3.0之后KRaft模式逐步替代ZooKeeper到3.5版本后ZooKeeper已经不是必选项。搭建时直接选KRaft模式少维护一个ZooKeeper集群也就少一层的故障点。三台机器组成的集群每台机器的配置文件 server.properties 核心参数如下process.rolesbroker,controller node.id1 controller.quorum.voters1192.168.1.11:9093,2192.168.1.12:9093,3192.168.1.13:9093 listenersPLAINTEXT://192.168.1.11:9092 advertised.listenersPLAINTEXT://192.168.1.11:9092 log.dirs/data/kafka-logs num.partitions6 default.replication.factor3 log.retention.hours24这个配置里最关键的是controller.quorum.voters它定义了集群里哪几个节点参与控制器选举三台机器都要一致。advertised.listeners配置的是对外广播的地址必须写成Spark节点能访问到的IP千万不要写localhost否则Spark消费端连过来会拿到一个连不通的地址。log.retention.hours设成24小时这是流处理场景的常见设置——实时分析不依赖超过一天的历史消息留太长徒增磁盘占用留太短一旦Spark任务故障超过这个时间窗口offset已经过期数据追不回来。启动顺序也有讲究第一台节点先启动等它初始化完成再启动另外两台。KRaft模式下首次启动需要执行kafka-storage.sh format命令格式化存储目录这个步骤在每台机器上都要做但只能在首次启动前执行一次重复执行会清空已有数据。3.3 Spark集群搭建Standalone模式足够Spark的集群模式有多套Standalone、YARN、Kubernetes。做这套系统直接上Standalone模式理由很简单一是部署链路短三台机器装好JDK、解压Spark、配好 slaves 文件就能启动二是不用去处理YARN队列和资源调度的概念把全部资源交给Spark自己管理。YARN模式以后有需要再迁切换成本不高。spark-env.sh 里的关键配置如下export JAVA_HOME/usr/lib/jvm/java-1.8.0 export SPARK_MASTER_HOST192.168.1.10 export SPARK_WORKER_CORES8 export SPARK_WORKER_MEMORY16g export SPARK_DRIVER_MEMORY4g注意SPARK_WORKER_MEMORY不要贪大不要把机器所有内存都分给Spark要预留出系统本身和Kafka进程需要的内存。我曾经把一台16G内存的机器给Spark分了14G结果Kafka broker频繁OOM查了半天才发现是内存打架了。另外JDK版本这里用的是8Spark 3.x用JDK 8完全兼容不要盲目上JDK 17某些版本的Spark在JDK 17下会出现反射相关的异常排查起来相当费时。3.4 连接字符串与参数初值Spark和Kafka之间的连接在Structured Streaming里是通过options传入的这里给出一个经过验证的初始配置val df spark.readStream .format(kafka) .option(kafka.bootstrap.servers, 192.168.1.11:9092,192.168.1.12:9092,192.168.1.13:9092) .option(subscribe, sensor_temperature,sensor_motion,device_switch) .option(startingOffsets, earliest) .option(failOnDataLoss, false) .option(maxOffsetsPerTrigger, 50000) .load()startingOffsets设为earliest表示从最早的未消费消息开始读适合首次上线想追历史数据的场景但如果Kafka topic里积压了大量测试数据第一次启动时会一下子读完所有消息需要几分钟甚至更久。failOnDataLoss设为false非常重要否则Kafka删除过期消息后Spark作业会因为找不到offset而报错停掉。maxOffsetsPerTrigger是流处理的节流阀控制每个微批最多处理多少条消息防止积压数据瞬间打爆下游存储。这些参数初值在生产环境跑几天后再根据实际延迟调整。4. 核心代码实现模拟设备数据写入Kafka并用Spark Structured Streaming做实时分析4.1 模拟设备端数据生产用量小但要真实没有真实智能家居设备的时候数据模拟器就是整个系统的发动机。模拟器要模拟出真实设备的行为特征而不是均匀地每秒发一条。真实设备上报有突发性、有周期规律也有数据跳跃。下面这个Python脚本模拟了一个温湿度传感器集群import json import random import time from kafka import KafkaProducer producer KafkaProducer( bootstrap_servers[192.168.1.11:9092, 192.168.1.12:9092], value_serializerlambda v: json.dumps(v).encode(utf-8) ) device_ids [fsensor_th_{i} for i in range(1, 201)] while True: # 每轮模拟一部分设备同时上报模拟设备分组上报的突发特征 active_batch random.sample(device_ids, krandom.randint(20, 60)) for device_id in active_batch: # 温度在22~28度之间波动偶尔出现设备故障导致的离谱读数 if random.random() 0.005: temperature round(random.uniform(70, 90), 1) else: temperature round(random.uniform(22, 28) random.uniform(-0.5, 0.5), 1) humidity round(random.uniform(40, 70), 1) timestamp_ms int(time.time() * 1000) message { deviceId: device_id, timestamp: timestamp_ms, type: temperature_humidity, data: {temperature: temperature, humidity: humidity} } producer.send(sensor_temperature, keydevice_id.encode(utf-8), valuemessage) # 模拟上报的突发每次批量上报后先短眠再随机休眠 time.sleep(random.uniform(0.3, 1.5))这个模拟器里有两个细节是经过考量的。第一每次随机选20到60个设备同时上报这个设定模拟了网关周期性聚集上报的真实行为Kafka接收到的消息不是平滑的而是脉冲式的。第二温度读数有0.5%的概率跳到70到90度这是在模拟设备故障或电磁干扰产生的脏数据——后面Spark的清洗逻辑会处理这类异常值模拟器必须喂一些脏数据清洗代码才有用武之地。producer.send指定了key这是为了让同一个设备的消息都进入Kafka的同一个partition这样Spark端做窗口聚合时同一设备的消息在同一个分区里天然有序。如果这里不设置keyKafka会按轮询策略把消息散到各个分区到Spark处理时乱序问题会很突出。4.2 Spark Structured Streaming消费与清洗Spark端的核心作业用Scala编写分为消费、清洗、聚合、写出四层结构。最基础的消费和过滤逻辑如下import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ val schema new StructType() .add(deviceId, StringType) .add(timestamp, LongType) .add(type, StringType) .add(data, new StructType() .add(temperature, DoubleType) .add(humidity, DoubleType)) val rawDF spark.readStream .format(kafka) .option(kafka.bootstrap.servers, 192.168.1.11:9092,192.168.1.12:9092) .option(subscribe, sensor_temperature) .option(startingOffsets, earliest) .option(maxOffsetsPerTrigger, 10000) .load() .selectExpr(CAST(value AS STRING) as json_str) val parsedDF rawDF .select(from_json(col(json_str), schema).as(msg)) .select(msg.deviceId, msg.timestamp, msg.data.temperature, msg.data.humidity) .withColumn(event_time, (col(timestamp) / 1000).cast(TimestampType)) // 清洗过滤掉物理上不可能出现的温湿度读数 val cleanDF parsedDF .filter(temperature BETWEEN -20 AND 60) .filter(humidity BETWEEN 0 AND 100)这段代码里from_json是结构解析的核心它需要一个预定义的schema作为参数schema定义必须和模拟器里写的JSON字段完全一致错一个字段名整个解析就会变成null。这里特别值得注意的一个坑模拟器里timestamp是毫秒数解析出来是LongType但流处理的时间窗口需要的是TimestampType所以要除以1000之后再cast成Timestamp。这个步骤漏掉的话后面做窗口聚合时时间维度全错。清洗阶段的filter条件看着简单但它决定了下游分析的质量。温度过滤范围定在-20度和60度之间这覆盖了智能家居场景所有合理可能——北方冬天暖气房、南方夏天、设备装在某些散热位置——同时又能把模拟器中那些70到90度的故障读数拦截掉。实际项目中清洗规则要比这个复杂但从最简单可靠的物理范围开始永远是对的。4.3 窗口聚合统计每分钟温度均值和异常计数清洗之后的数据要落到具体业务指标上。智能家居系统最常见的两个需求是实时掌握每个房间的温度趋势、及时发现设备异常上报。窗口聚合是实现这两个指标的最直接手段。// 1分钟滑动窗口每30秒触发一次统计 val windowedStats cleanDF .withWatermark(event_time, 30 seconds) .groupBy( col(deviceId), window(col(event_time), 1 minute, 30 seconds) ) .agg( avg(temperature).as(avg_temperature), max(temperature).as(max_temperature), count(when(col(temperature) 45, true)).as(abnormal_count) )withWatermark是这里最关键的参数。它声明了一个30秒的水位线含义是允许事件时间晚于当前最大事件时间30秒以内的数据参与计算更晚的数据直接丢弃。为什么是30秒而不是更长因为模拟器里消息的Kafka延迟在毫秒级网络抖动也不会超过几秒30秒足够覆盖正常的迟延范围同时又能把严重迟到的那批数据挡在外面。如果把水位线设成10分钟那么窗口要一直保持10分钟才能关闭内存里会积压大量中间状态而且下游看到的结果会延迟很长时间。window(col(event_time), 1 minute, 30 seconds)这个窗口函数的参数含义也要说清楚第一个参数是时间列第二个参数是窗口长度第三个参数是滑动步长。这个配置下每分钟的数据会被切分成两个重叠的30秒窗口计算每30秒出一轮结果数据重复计算两遍。如果改成长度5分钟、步长1分钟那就是标准的滚动窗口式输出。窗口到底设多大核心取决于业务上“快”到什么程度算快——做实时大屏30秒以下比较合理做报表分析窗口可以放大到5分钟甚至更长。4.4 结果写出Redis写入与MySQL批量落库聚合结果算出来后要同时解决两个读者实时大屏需要快速拉取最新数据报表系统需要明细数据落库。这两者的写入路径完全不同。// 写入Redis按设备ID存最近一条聚合结果供大屏实时拉取 val redisWriter windowedStats .writeStream .foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, _: Long) batchDF.collect().foreach { row val deviceId row.getString(0) val avgTemp row.getDouble(2) redisClient.set(srealtime:temp:$deviceId, avgTemp.toString) redisClient.expire(srealtime:temp:$deviceId, 90) } } .outputMode(update) .start() // 写入MySQL窗口结果追加到明细表供离线报表使用 val mysqlWriter windowedStats .writeStream .foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, _: Long) batchDF.write .mode(append) .jdbc( jdbc:mysql://192.168.1.20:3306/smart_home, realtime_device_stats, Properties() ) } .outputMode(append) .start()两个写出都用了foreachBatch模式这是Structured Streaming里最灵活的写出方式——它把流处理里的每个微批变成一个独立的DataFrame然后我可以对这个DataFrame调用任何批处理API。选择foreachBatch而不是默认的foreach或writeStream.format原因很实际foreachBatch可以复用批处理的连接池和写入端到端JDBC写入天然支持批量插入性能比逐条foreach好一个数量级。写入Redis时设置了90秒过期这个设计意在大屏上自动清理掉已经不再上报的设备——如果某台设备停止上报90秒后它的缓存键自动消失前端设备列表里就不会再显示这台离线设备。MySQL写入用mode(append)只追加不更新因为明细数据是按时间累积的历史结果不应该被改动。4.5 完整启动作业时的三板斧Spark流任务跑起来之后第一件事是看日志、看UI、看数据产出这三件事有固定的检查顺序。启动命令如下spark-submit \ --master spark://192.168.1.10:7077 \ --class com.smarthome.RTDeviceAnalyzer \ --executor-memory 6g \ --executor-cores 4 \ smart-home-analytics.jar提交后先看Spark UI的Streaming标签页Web UI地址是http://192.168.1.10:8080重点看两个指标Input Rate每秒接收的消息数和Scheduling Delay调度延迟。如果Scheduling Delay持续上升说明Spark的处理速度跟不上Kafka的消息生产速度这时需要看maxOffsetsPerTrigger是否限制得太死或者executor的内存和并行度不够。再查Redis里的realtime:temp:*key是否有持续更新的值如果Redis里有值但Kafka消费速率却为0基本可以断定是Spark作业没有启动成功回去翻stderr日志。5. 避坑与排查Spark与Kafka联调中常见的5个典型问题5.1 消费不到数据但Kafka里明明有生产消息现象Kafka用命令行工具能看到生产消息数量在增长但Spark作业处理速率为零UI显示的Input Rate始终是0。原因85%的情况是kafka.bootstrap.servers写的是localhost或者写的是Kafka节点在/etc/hosts里的内网主机名而Spark executor所在的机器解析不了这个主机名。Kafka broker会把advertised.listeners里配置的地址返回给消费者如果消费者连不上这个地址就会反复拉取元数据失败但不会报错退出表现就是静默地消费不到数据。解决逐项排查。先登上Spark executor所在的机器依次执行telnet 192.168.1.11 9092、kafka-console-consumer.sh --bootstrap-server 192.168.1.11:9092 --topic sensor_temperature --from-beginning验证网络通不通、命令行能不能消费。如果命令行消费正常而Spark消费不到那再检查Spark默认配置里的spark.executor.extraJavaOptions是否设置了不兼容的Kafka客户端参数比如security.protocol或sasl.*相关的配置残留。5.2 流作业跑着跑着自己停了没有报错现象Spark UI里作业状态变成FINISHED或DEAD日志里只看到一行INFO级别的“Shutting down Query”没有任何Exception。原因这是Structured Streaming里最隐蔽的问题之一——代码里某个分支调用了df.show()或df.collect()并且指定了每隔一段时间触发一次而这些Action操作默认会停掉流查询。很多新手在调试时为了看一眼中间结果在流处理逻辑里加了一行parsedDF.show(10)调试完了忘记删结果作业每次跑到某个时刻就被这条语句中断。解决全局搜索代码里的.show()和.collect()凡是出现在writeStream链之外的都删掉。如果确实需要在流处理过程中观察中间数据正确做法是单独写一个临时的writeStream.format(console).start()看完后手动stop()或者把中间结果写入一个临时Kafka topic再消费查看。5.3 窗口聚合结果延迟了5到10分钟才出现象Redis里的聚合结果刷新频率远低于窗口配置的30秒有时甚至5分钟才更新一次。原因这是三个参数叠加导致的。第一上游Kafka的topic分区数设为默认的单个分区Spark对单个分区只能起一个消费线程吞吐受限。第二maxOffsetsPerTrigger设置过小每个微批只消费很少的消息处理进度慢。第三withWatermark的滞后时间设置太长导致窗口迟迟不触发末次计算。解决把Kafka topic分区数至少扩到12个同时调大maxOffsetsPerTrigger。检查方式是在Spark UI的Streaming页签看Batch Processing Time如果每个批处理耗时超过触发间隔并行度一定不够。分区数的扩大方法是新建topic时指定分区数不要在旧topic上alter partitions——旧topic分区扩容会把key的分区重排造成一段时间的数据乱序。5.4 模拟器端偶发报错 ProducerSendException现象模拟器跑一段时间后控制台偶尔出现“Batch Expired”或“ProducerSendException”消息发送失败部分数据丢失。原因模拟器生产的消息突发性太强。time.sleep(random.uniform(0.3, 1.5))这个短眠结束后一次性发送几十条消息而Kafka的batch.size默认16KB、linger.ms默认0意味着每个批次消息只要达到16KB就立即独立发送瞬时产生大量网络请求而Python的KafkaProducer默认的max_in_flight_requests_per_connection为5超过限制的请求会被拒绝。解决修改模拟器在创建KafkaProducer时的参数把linger_ms调大到500batch_size调大到65536并设置retries3。linger_ms增大后消息会在内存中多等500毫秒凑成一个更大的批次再发送网络请求次数大幅减少。同时给生产者加一个回调函数在发送失败时打印错误信息并重试确保测试期间不丢数据def on_send_failure(exc): if exc is not None: print(fSend failed: {exc})5.5 任务重启后从重复数据开始消费现象Spark作业重启后发现MySQL里出现了重复的窗口聚合结果同一个设备同一个时间窗的数据被写入了两次。原因这其实是正确的行为不是bug。Structured Streaming默认提供端到端的至少一次at-least-once保证Spark自动把offset信息存储在checkpoint目录里作业重启后从checkpoint恢复消费位置不会重复消费Kafka消息。但foreachBatch写出到MySQL的这个环节如果MySQL写入执行成功但Spark还没提交offset此时作业崩溃重启后就会重新处理这个批次造成重复写入。解决在写入MySQL时做幂等处理。最稳妥的做法是在目标表上增加唯一约束比如(device_id, window_start)联合唯一键插入时用INSERT ... ON DUPLICATE KEY UPDATE或者先DELETE该时间窗再INSERT。在foreachBatch内部写一个batchDF.write.mode(overwrite).jdbc(...)配合MySQL表的唯一索引即可把重复写入收敛成一次更新。6. 验证与进阶如何判断这套系统真的跑通了以及从哪里入手优化系统跑起来之后不要只看Spark UI的Input Rate就认为万事大吉那只是Kafka数据的流入速度不是业务结果的正确性。第一步验证数据质量从Redis里随机抽几台设备的温度和MySQL里的窗口明细对比两者不能有数量级差异。第二步验证延迟在模拟器里打印一条消息的生产时间然后在Redis里看到这条设备聚合结果的时间两者差值就是端到端延迟正常应该在3到5秒以内。第三步模拟故障手动杀掉Spark作业进程等两分钟再重启检查这段时间Kafka里累积的消息是否全部被处理、MySQL里有没有出现主键冲突的报错。进阶优化的第一优先级是调整并行度。当前的设计里Kafka的partition数决定了Spark的并发读取能力如果发现CPU利用率一直很低但数据延迟偏高大概率是partitions数量不足。把Kafka的topic分区从6扩到24同时把Spark的spark.sql.shuffle.partitions从默认的200调整到分区数的两倍能明显感觉到吞吐量提升。第二优先级是把设备分组上报做降噪目前模拟器的异常读数直接过滤掉了实际场景里更合理的是将异常值标记而不是丢弃保留一个is_abnormal字段后续离线分析能用它来定位故障设备。这个系统做下来最大的收获是理解了流处理架构里“缓冲层和计算层解耦”的价值Kafka扛突刺、Spark管计算两层各自独立演进线上排查问题也只需要盯着一个方向。我从最初不停在模拟器、Kafka、Spark之间来回折腾到现在基本能根据UI指标一分钟定位问题所在中间踩过不少坑也总结了一套自己的调试习惯——先确认数据是否真的进了Kafka再确认Spark是否真的消费了最后确认结果是否真的写对了沿着这条线走90%的问题都能快速收敛。希望帮到你。本文还有配套的精品资源点击获取
返回列表