【问题标题】:Getting Hadoop OutputFormat RunTimeException while running Apache Spark Kafka Stream运行 Apache Spark Kafka Stream 时获取 Hadoop OutputFormat RunTimeException
【发布时间】:2016-11-25 00:04:28
【问题描述】:

我正在运行一个程序,该程序使用 Apache Spark 从 Apache Kafka 集群获取数据并将数据放入 Hadoop 文件中。我的程序如下:

public final class SparkKafkaConsumer {
    public static void main(String[] args) {
        SparkConf sparkConf = new SparkConf().setAppName("JavaKafkaWordCount");
        JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(2000));
        Map<String, Integer> topicMap = new HashMap<String, Integer>();
        String[] topics = "Topic1, Topic2, Topic3".split(",");
        for (String topic: topics) {
            topicMap.put(topic, 3);
        }
        JavaPairReceiverInputDStream<String, String> messages =
                KafkaUtils.createStream(jssc, "kafka.test.com:2181", "NameConsumer", topicMap);
        JavaDStream<String> lines = messages.map(new Function<Tuple2<String, String>, String>() {
            public String call(Tuple2<String, String> tuple2) {
                return tuple2._2();
            }
        });
        JavaDStream<String> words = lines.flatMap(new FlatMapFunction<String, String>() {
            public Iterable<String> call(String x) {
                return Lists.newArrayList(",".split(x));
            }
        });
        JavaPairDStream<String, Integer> wordCounts = words.mapToPair(
                new PairFunction<String, String, Integer>() {
                    public Tuple2<String, Integer> call(String s) {
                        return new Tuple2<String, Integer>(s, 1);
                    }
                }).reduceByKey(new Function2<Integer, Integer, Integer>() {
                    public Integer call(Integer i1, Integer i2) {
                        return i1 + i2;
                    }
                });
        wordCounts.print();
        wordCounts.saveAsHadoopFiles("hdfs://localhost:8020/user/spark/stream/", "txt");
        jssc.start();
        jssc.awaitTermination();
    }
}

我正在使用这个命令提交申请:C:\spark-1.6.2-bin-hadoop2.6\bin\spark-submit --packages org.apache.spark:spark-streaming-kafka_2.10:1.6.2 --class "SparkKafkaConsumer" --master local[4] target\simple-project-1.0.jar

我收到此错误:java.lang.RuntimeException: class scala.runtime.Nothing$ not org.apache.hadoop.mapred.OutputFormat at org.apache.hadoop.conf.Configuration.setClass(Configuration.java:2148)

是什么导致了这个错误,我该如何解决?

【问题讨论】:

  • 这看起来像是 Spark stackoverflow.com/questions/29007085/… 中的一个问题 ..
  • 你可以试试saveAsHadoopFiles("hdfs://localhost:8020/user/spark/stream/", "txt", Text.class, IntWritable.class, TextOutputFormat.class)吗?
  • @Hawknight Text.classTextOutputFormat.class的完整包是什么?
  • Text 来自org.apache.hadoop.ioTextOutputFormat 来自org.apache.hadoop.mapred
  • @Hawknight 解决了它。谢谢!您应该将此添加为答案。

标签: java scala hadoop apache-spark apache-kafka


【解决方案1】:

我同意该错误并不是真正令人回味的,但通常最好在任何 saveAsHadoopFile 方法中指定要输出的数据格式,以保护自己免受此类异常的影响。

这是文档中特定方法的原型:

saveAsHadoopFiles(java.lang.String prefix, java.lang.String suffix, java.lang.Class<?> keyClass, java.lang.Class<?> valueClass, java.lang.Class<F> outputFormatClass)

在您的示例中,这将对应于:

wordCounts.saveAsHadoopFiles("hdfs://localhost:8020/user/spark/stream/", "txt", Text.class, IntWritable.class, TextOutputFormat.class)

根据您的wordCounts PairDStream 的格式,我选择了Text,因为键的类型为String,而IntWritable,因为与键关联的值的类型为Integer

如果您只需要基本的纯文本文件,请使用TextOutputFormat,但您可以查看FileOutputFormat 的子类以获得更多输出选项。

同样被问到,Text 类来自org.apache.hadoop.io 包,TextOutputFormat 来自org.apache.hadoop.mapred 包。

【讨论】:

    【解决方案2】:

    为了完整性(@Jonathan 给出了正确答案)

    import org.apache.hadoop.io.IntWritable;
    import org.apache.hadoop.io.Text;
    import org.apache.hadoop.mapred.TextOutputFormat;
    
    ...
    wordCounts.saveAsHadoopFiles("hdfs://localhost:8020/user/spark/stream/", "txt", Text.class, IntWritable.class, TextOutputFormat.class)
    

    【讨论】:

      猜你喜欢
      • 2016-11-01
      • 2016-06-26
      • 1970-01-01
      • 2012-03-03
      • 2019-04-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-01-19
      相关资源
      最近更新 更多