Spark-Submit执行Scala Kafka流程序遇scala.MatchError问题求助
Scala Spark Kafka流程序spark-submit时出现scala.MatchError异常排查
错误日志
Exception in thread "main" scala.MatchError: [Ljava.lang.String;@368d5c00 (of class [Ljava.lang.String;) at com.huawei.bigdata.spark.examples.KafkaToOBS$.createContext(KafkaToOBS-test1.scala:30) at com.huawei.bigdata.spark.examples.KafkaToOBS$.main(KafkaToOBS-test1.scala:22)
程序代码
package com.huawei.bigdata.spark.examples import org.apache.hadoop.mapred.lib.MultipleTextOutputFormat import org.apache.hadoop.conf.Configuration import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010._ import java.time.{Instant, LocalDateTime, ZoneId} import java.time.format.DateTimeFormatter /** * Consumes messages from one or more topics in Kafka. * <checkPointDir> is the Spark Streaming checkpoint directory. * <brokers> is for bootstrapping and the producer will only use it for getting metadata * <topics> is a list of one or more kafka topics to consume from * <batchTime> is the Spark Streaming batch duration in seconds. */ object KafkaToOBS { def main(args: Array[String]) { val ssc = createContext(args) // (KafkaToOBS-test1.scala:22) //The Streaming system starts. ssc.start() ssc.awaitTermination() } def createContext(args : Array[String]) : StreamingContext = { val Array(checkPointDir, brokers, topics, batchTime, groupId, path) = args //(KafkaToOBS-test1.scala:30) // Create a Streaming startup environment. val sparkConf = new SparkConf().setAppName("KafkaToOBS") val ssc = new StreamingContext(sparkConf, Seconds(batchTime.toLong)) //Configure the CheckPoint directory for the Streaming. //This parameter is mandatory because of existence of the window concept. ssc.checkpoint(checkPointDir) // Get the list of topic used by kafka val topicArr = topics.split(",") val topicSet = topicArr.toSet val kafkaParams = Map[String, String]( "bootstrap.servers" -> brokers, "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "group.id" -> groupId, "auto.offset.reset" -> "earliest" ); val locationStrategy = LocationStrategies.PreferConsistent val consumerStrategy = ConsumerStrategies.Subscribe[String, String](topicSet, kafkaParams) // Create direct kafka stream with brokers and topics // Receive data from the Kafka and generate the corresponding DStream val stream = KafkaUtils.createDirectStream[String, String](ssc, locationStrategy, consumerStrategy) val tranData = stream.map(x=>(x.value(),"")) tranData.repartition(1).foreachRDD(result=>{ result.saveAsHadoopFile(path, classOf[String], classOf[String],classOf[RDDMultipleTextOutputFormat]) }) ssc } class RDDMultipleTextOutputFormat extends MultipleTextOutputFormat[Any, Any]{ def coverTimeStampToString(time:Long): String ={ val format = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss") format.format(LocalDateTime.ofInstant(Instant.ofEpochMilli(time),ZoneId.systemDefault())) } override def generateFileNameForKeyValue(key: Any, value: Any, name: String):String ={ val timeStamp = System.currentTimeMillis() coverTimeStampToString(timeStamp) val ymd=getDay(timeStamp) val hour=getHour(timeStamp) val service_date="day="+ymd+"/"+"hour="+hour+"/"+name+"_"+getMinute(timeStamp)//Writing path: test\streaming\day=2020-07-10\hour=17\part-00000_26 service_date } def getHour(time:Long): String ={ coverTimeStampToString(time).substring(11,13) } def getDay(time:Long): String ={ coverTimeStampToString(time).substring(0,10) } def getMinute(time:Long): String ={ coverTimeStampToString(time).substring(14,16) } } }
spark-submit命令
spark-submit --master yarn --deploy-mode client --num-executors 3 --class com.huawei.bigdata.spark.examples.KafkaToOBS --jars $(files= /opt/client/Spark2x/spark/jars/*.jar),SparkStreamingKafka010Example-1.0.jar hdfs://hacluster/tmp/ 192.xxx.x.xx:9092 DemoPOC 10 groupaaa obs://sg-demo-poc/KafkaToOBS/
排查与解决思路
1. 定位参数传递问题
报错出现在代码第30行的数组匹配处,本质是传入的参数数量不等于代码预期的6个。问题根源在spark-submit命令的--jars参数写法错误:
$(files= /opt/client/Spark2x/spark/jars/*.jar)是无效的shell语法,不会生成jar列表,反而会将这个字符串作为一个jar路径传入,导致后续的位置参数被整体偏移,实际传入程序的参数数量和内容不符合预期。
2. 修正spark-submit命令
修改--jars部分的写法,避免错误传递参数:
- 如果不需要额外引入Spark自带jar,直接简化命令:
spark-submit --master yarn --deploy-mode client --num-executors 3 \ --class com.huawei.bigdata.spark.examples.KafkaToOBS \ --jars SparkStreamingKafka010Example-1.0.jar \ hdfs://hacluster/tmp/ 192.xxx.x.xx:9092 DemoPOC 10 groupaaa obs://sg-demo-poc/KafkaToOBS/
- 如果确实需要添加额外jar,正确生成逗号分隔的列表:
spark-submit --master yarn --deploy-mode client --num-executors 3 \ --class com.huawei.bigdata.spark.examples.KafkaToOBS \ --jars "$(ls /opt/client/Spark2x/spark/jars/*.jar | tr '\n' ',' | sed 's/,$//'),SparkStreamingKafka010Example-1.0.jar" \ hdfs://hacluster/tmp/ 192.xxx.x.xx:9092 DemoPOC 10 groupaaa obs://sg-demo-poc/KafkaToOBS/
3. 增强代码的参数校验
在代码中添加参数数量校验,避免直接匹配数组抛出模糊异常:
def createContext(args : Array[String]) : StreamingContext = { if (args.length != 6) { throw new IllegalArgumentException(s"预期6个参数,实际收到${args.length}个:${args.mkString(",")}") } val Array(checkPointDir, brokers, topics, batchTime, groupId, path) = args // 后续代码... }
4. 打印参数详情排查
在main方法开头添加打印语句,确认实际传入的参数:
def main(args: Array[String]) { println(s"参数数量:${args.length}") println(s"参数列表:${args.mkString("[", ",", "]")}") val ssc = createContext(args) ssc.start() ssc.awaitTermination() }
通过日志可以直观看到参数是否正确传递,快速定位问题。
内容的提问来源于stack exchange,提问作者user2778168
相关产品推荐
相关产品推荐

