You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 16:25:24