-1
這是一個字數映射縮減作業。我有我自己的InputFormat。Reducer任務未在我的MapReduce作業中調用
個JobExecutor:
val job = new Job(new Configuration())
job.setMapperClass(classOf[CountMapper])
job.setReducerClass(classOf[CountReducer])
job.setJobName("tarun-test-1")
job.setInputFormatClass(classOf[MyInputFormat])
FileInputFormat.setInputPaths(job, new Path(args(0)))
FileOutputFormat.setOutputPath(job, new Path(args(1)))
job.setOutputKeyClass(classOf[Text])
job.setOutputValueClass(classOf[LongWritable])
job.setNumReduceTasks(1)
println("status: " + job.waitForCompletion(true))
映射:
class CountMapper extends Mapper[LongWritable, Text, Text, LongWritable] {
private val valueOut = new LongWritable(1L)
override def map(k: LongWritable, v: Text, context: Mapper[LongWritable, Text, Text, LongWritable]#Context): Unit = {
val str = v.toString
str.split(",").foreach(word => {
val keyOut = new Text(word.toLowerCase.trim)
context.write(keyOut, valueOut)
})
}
}
減速機:
class CountReducer extends Reducer[Text, LongWritable, Text, LongWritable] {
override def reduce(k: Text, values: Iterable[LongWritable], context: Reducer[Text, LongWritable, Text, LongWritable]#Context): Unit = {
println("Inside reduce method..")
val valItr = values.iterator()
var sum = 0L
while (valItr.hasNext) {
sum = sum + valItr.next().get()
}
context.write(k, new LongWritable(sum))
println("done reducing.")
}
}
映射被調用,RecordReader正在讀正確基於日誌分裂。但是,reducer沒有被調用。
你是什麼意思你有你自己的InputFormat?它在哪裏?你的意思是減少沒有被調用?你怎麼知道?任何輸入/輸出?計數器?錯誤?日誌? – vefthym
MyInputFormat是我自己的InputFormat。 InputFormat按預期工作,我看到映射器的輸入(鍵,值)正在被RecordReader正確讀取。我將日誌記錄添加到Map任務,並按預期記錄事件。但是,減少日誌不會打印並且最終狀態爲false。 –