Why does implicit conversions for Writable doesn't work

Viewed 946

SparkContext defines a couple of implicit conversions between Writable and their primitive types, like LongWritable <-> Long, Text <-> String.

  • TEST CASE 1:

I am using the following code to combine small files

  @Test
  def  testCombineSmallFiles(): Unit = {
    val path = "file:///d:/logs"
    val rdd = sc.newAPIHadoopFile[LongWritable,Text, CombineTextInputFormat](path)
    println(s"rdd partition number is ${rdd.partitions.length}")
    println(s"lines is :${rdd.count()}")
  }

The above code works well, but If I use the following line to get the rdd, it will result in a compiling error:

val rdd = sc.newAPIHadoopFile[Long,String, CombineTextInputFormat](path)

It looks that the implicit conversion doesn't take effect. I would like to know what's wrong here and why it isn't working.

  • TEST CASE 2:

With following code that is using sequenceFile, the implicit conversion looks working(Text is converted to String, and IntWritable is converted to Int)

 @Test
  def testReadWriteSequenceFile(): Unit = {
    val data = List(("A", 1), ("B", 2), ("C", 3))
    val outputDir = Utils.getOutputDir()
    sc.parallelize(data).saveAsSequenceFile(outputDir)
    //implicit conversion works for the SparkContext#sequenceFile method
    val rdd = sc.sequenceFile(outputDir + "/part-00000", classOf[String], classOf[Int])
    rdd.foreach(println)
  }

Compare these two test cases, I didn't see the key difference that make's one work,and the other doesn't work.

  • NOTE:

The SparkContext#sequenceFile method that I am using in the TEST CASE 2 is:

  def sequenceFile[K, V](
      path: String,
      keyClass: Class[K],
      valueClass: Class[V]): RDD[(K, V)] = withScope {
    assertNotStopped()
    sequenceFile(path, keyClass, valueClass, defaultMinPartitions)
  }

In the sequenceFile method,it is calling another sequenceFile method, which is calling hadoopFile method to read the data

  def sequenceFile[K, V](path: String,
      keyClass: Class[K],
      valueClass: Class[V],
      minPartitions: Int
      ): RDD[(K, V)] = withScope {
    assertNotStopped()
    val inputFormatClass = classOf[SequenceFileInputFormat[K, V]]
    hadoopFile(path, inputFormatClass, keyClass, valueClass, minPartitions)
  }
1 Answers
Related