使用KafkaUtils.createDirectStream时Hive无更新、Spark Streaming无日志输出求助
解决Spark Streaming对接Kafka无输出、无数据写入的问题
首先看你的代码,我发现了最关键的问题:你的StreamingContext没有启动,也没有让程序保持运行的逻辑,这直接导致整个流处理的核心逻辑(foreachRDD里的代码)完全没被执行——所以你只会看到JobScheduler添加作业的日志,但实际没有任何数据处理动作,自然也就没有stdout输出、Hive表也不会更新。
核心修复步骤
你必须在代码末尾添加这两行关键代码:
ssc.start() // 启动StreamingContext,开始接收Kafka数据并处理 ssc.awaitTermination() // 让程序保持运行状态,等待手动终止或异常退出
其他需要排查的细节
除了这个核心问题,还可以检查以下几点:
- Kafka主题是否有数据:先用
kafka-console-consumer.sh工具测试topic主题里是否有数据,如果主题本身是空的,foreachRDD的处理逻辑也不会触发 - Kafka连接参数是否正确:确认
metadata.broker.list的IP和端口能正常访问,group.id配置无误;auto.offset.reset设为smallest是从最旧数据开始消费,但如果这个消费者组之前已经消费过该主题的全部数据,也会出现无数据可处理的情况 - 日志输出位置:如果是
yarn-client模式,stdout输出会在客户端控制台;如果是yarn-cluster模式,需要通过yarn logs -applicationId <你的appId>查看完整日志;本地模式下则直接输出在控制台 - Hive表写入权限:确保Spark程序有写入
test.rt_test表的权限,表的存储路径配置正确;另外SaveMode.Overwrite是覆盖写入,只有当当前批次有数据时才会更新表内容
修复后的完整代码示例
val sparkConf = new SparkConf().setAppName("DirectKafka") sparkConf.setMaster("yarn-client") // 也可以尝试local[*] // sparkConf.set("spark.streaming.kafka.maxRatePerPartition", "5") // sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") val ssc = new StreamingContext(sparkConf, Seconds(2)) val topicsSet = Set("topic") val kafkaParams = Map[String, String]( "metadata.broker.list" -> "ip1:9092,ip2:9092", "group.id" -> "test0107", "auto.offset.reset" -> "smallest" ) // 补全createDirectStream的完整参数(根据你的Spark版本调整泛型和解码器) val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder]( ssc, kafkaParams, topicsSet ) messages.foreachRDD{rdd => val spark = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate() import spark.implicits._ val df = spark.read.json(rdd.map(_._2)) println(df.count()) // 用println替代print,确保输出换行易读 df.createOrReplaceTempView("words") val wc = spark.sql("select count(1) as id from words ") wc.write.mode(SaveMode.Overwrite).insertInto("test.rt_test") wc.show() } // 启动流程序并保持运行 ssc.start() ssc.awaitTermination()
内容的提问来源于stack exchange,提问作者Autumn
相关产品推荐
相关产品推荐

