【问题标题】:how to store processed data from hdfs using mapReduce in mongoDB as output如何使用 mongoDB 中的 mapReduce 从 hdfs 存储处理过的数据作为输出
【发布时间】:2015-09-11 03:57:42
【问题描述】:

我有一个 mapreduce 应用程序,它处理来自 HDFS 的数据并将输出数据存储在 HDFS 中

但是,现在我需要将输出数据存储在 mongodb 中,以便将其存储到 HDFS 中

谁能告诉我怎么做?

谢谢

映射器类

package com.mapReduce;

import java.io.IOException;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

public class FMapper extends Mapper<LongWritable, Text, Text, Text> {
    private String pART;
    private String actual;
    private String fdate;
    public void map(LongWritable ikey, Text ivalue, Context context) throws IOException, InterruptedException {
        String tempString = ivalue.toString();
        String[] data = tempString.split(",");
        pART=data[1];
        try{
            fdate=convertyymmdd(data[0]);
            /**IF ACTUAL IS LAST HEADER
             * actual=data[2];
             * */
            actual=data[data.length-1];
            context.write(new Text(pART), new Text(fdate+","+actual+","+dynamicVariables(data)));
        }catch(ArrayIndexOutOfBoundsException ae){
            System.err.println(ae.getMessage());
        }

    }


    public static String convertyymmdd(String date){

        String dateInString=null;
        String data[] =date.split("/");
        String month=data[0];
        String day=data[1];
        String year=data[2];
        dateInString =year+"/"+month+"/"+day;
        System.out.println(dateInString);   
        return dateInString;
    }

    public static String dynamicVariables(String[] data){
        StringBuilder str=new StringBuilder();
        boolean isfirst=true; 
    /** IF ACTUAL IS LAST HEADER
     * for(int i=3;i<data.length;i++){ */
        for(int i=2;i<data.length-1;i++){

            if(isfirst){
                str.append(data[i]);
                isfirst=false;
            }
            else
            str.append(","+data[i]);
        }
        return str.toString();
        }

}

减速器类

package com.mapReduce;

import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;

import javax.faces.bean.ApplicationScoped;
import javax.faces.bean.ManagedBean;
import javax.faces.bean.ManagedProperty;

import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;

import com.ihub.bo.ForcastBO;
import com.ihub.service.ForcastService;
public class FReducer extends Reducer<Text, Text, Text, Text> {
    private String pART;
    private List<ForcastBO> list = null;
    private List<List<String>> listOfList = null;
    private List<String> vals = null;
    private static List<ForcastBO> forcastBos=new ArrayList<ForcastBO>();

    @Override
    public void reduce(Text _key, Iterable<Text> values, Context context) throws IOException, InterruptedException {
    // TODO Auto-generated method stub
        pART = _key.toString();
        // process values
        for (Text val : values) {
            String tempString = val.toString();
            String[] data = tempString.split(",");
            ForcastBO fb=new ForcastBO();
            fb.setPart(pART);
            fb.setDate(data[0]);
            fb.setActual(data[1]);
            fb.setW0(data[2]);
            fb.setW1(data[3]);
            fb.setW2(data[4]);
            fb.setW3(data[5]);
            fb.setW4(data[6]);
            fb.setW5(data[7]);
            fb.setW6(data[8]);
            fb.setW7(data[9]);
            try {
                list.add(fb);
            } catch (Exception ae) {
                System.out.println(ae.getStackTrace() + "****" + ae.getMessage() + "*****" + ae.getLocalizedMessage());
            }
        }   
    }

    @Override
    public void run(Context context) throws IOException, InterruptedException {
        setup(context);
        try {
          while (context.nextKey()) {

         listOfList = new ArrayList<List<String>>();
         list=new ArrayList<ForcastBO>();
            reduce(context.getCurrentKey(), context.getValues(), context);
            files_WE(listOfList, list, context);

          }

          }finally {
              cleanup(context);
            }
    }


    public void files_WE(List<List<String>> listOfList, List<ForcastBO> list, Context context) {

        Collections.sort(list);

            try {
                setData(listOfList, list);

                Collections.sort(listOfList, new Comparator<List<String>>() {
                    @Override
                    public int compare(List<String> o1, List<String> o2) {
                        return o1.get(0).compareTo(o2.get(0));
                    }
                });

                for (int i = listOfList.size() - 1; i > -1; i--) {
                    List<String> list1 = listOfList.get(i);
                    int k = 1;
                    for (int j = 3; j < list1.size(); j++) {
                        try {
                            list1.set(j, listOfList.get(i - k).get(j));
                        } catch (Exception ex) {
                            list1.set(j, null);
                        }
                        k++;
                    }

                }
            } catch (Exception e) {
                //e.getLocalizedMessage();
            }

            for(List<String> ls:listOfList){
                System.out.println(ls.get(0));
                ForcastBO forcastBO=new ForcastBO();
                try{
                    forcastBO.setPart(ls.get(0));
                    forcastBO.setDate(ls.get(1));
                    forcastBO.setActual(ls.get(2));
                    forcastBO.setW0(ls.get(3));
                    forcastBO.setW1(ls.get(4));
                    forcastBO.setW2(ls.get(5));
                    forcastBO.setW3(ls.get(6));
                    forcastBO.setW4(ls.get(7));
                    forcastBO.setW5(ls.get(8));
                    forcastBO.setW6(ls.get(9));
                    forcastBO.setW7(ls.get(10));
                    forcastBos.add(forcastBO);
                    }catch(Exception e){
                        forcastBos.add(forcastBO);
                    }
                try{
                    System.out.println(forcastBO);
                    //service.setForcastBOs(forcastBos);
            }catch(Exception e){
                System.out.println("FB::::"+e.getStackTrace());
            }
            }
    }





        public void setData(List<List<String>> listOfList, List<ForcastBO> list) {
            List<List<String>> temListOfList=new ArrayList<List<String>>();
            for (ForcastBO str : list) {
                vals = new ArrayList<String>();
                vals.add(str.getPart());
                vals.add(str.getDate());
                vals.add(str.getActual());
                vals.add(str.getW0());
                vals.add(str.getW1());
                vals.add(str.getW2());
                vals.add(str.getW3());
                vals.add(str.getW4());
                vals.add(str.getW5());
                vals.add(str.getW6());
                vals.add(str.getW7());
                temListOfList.add(vals);
            }


            Collections.sort(temListOfList, new Comparator<List<String>>() {
                @Override
                public int compare(List<String> o1, List<String> o2) {
                    return o1.get(1).compareTo(o2.get(1));
                }
            });

            for(List<String> ls:temListOfList){
                System.out.println(ls);
                listOfList.add(ls);
                }
        }

        public static List<ForcastBO> getForcastBos() {
            return forcastBos;
        }



    }

驱动类

package com.mapReduce;

import java.net.URI;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;


public class MRDriver {

    public static void main(String[] args)  throws Exception {
        Configuration conf = new Configuration();
        Job job = Job.getInstance(conf, "JobName");
        job.setJarByClass(MRDriver.class);
        // TODO: specify a mapper
        job.setMapperClass(FMapper.class);
        // TODO: specify a reducer
        job.setReducerClass(FReducer.class);

        // TODO: specify output types
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(Text.class);

        // TODO: delete temp file
        FileSystem hdfs = FileSystem.get(new URI("hdfs://localhost:9000"),
                conf); 
        Path workingDir=hdfs.getWorkingDirectory();

        Path newFolderPath= new Path("/sd1");
        newFolderPath=Path.mergePaths(workingDir, newFolderPath);
        if(hdfs.exists(newFolderPath))

        {
            hdfs.delete(newFolderPath); //Delete existing Directory

        }
        // TODO: specify input and output DIRECTORIES (not files)

        FileInputFormat.setInputPaths(job,new Path("hdfs://localhost:9000/Forcast/SampleData"));
        FileOutputFormat.setOutputPath(job, newFolderPath);

        if (!job.waitForCompletion(true))
            return;
    }
}

【问题讨论】:

  • 首先您需要代码从 HDFS“读取”,然后您需要一个 MongoDB 驱动程序并将您的“写入”代码写入 MongoDB,或者根据需要从您的“减速器”或最终阶段直接输出到 MongoDB .基本上为您的语言获取一个驱动程序(hadoop 确实支持几种不同模式,但也许您的意思是 Java),然后连接并编写。不过先学习驱动程序。
  • 您处理的数据是什么格式的?您始终可以在 reducer 中调用 MongoDB 客户端并将数据批量写入清理部分(例如)。如果您希望我们提供帮助,请提供更多细节。
  • 处理后的数据是LIST格式
  • 如果需要我可以添加代码
  • 从上面的mapreduce我将数据存储在hdfs中,但现在我需要mongodb中的数据

标签: mongodb hadoop


【解决方案1】:

基本上你需要改变“输出格式类”,你有几种方法:

  1. 使用 Hadoop 的 MongoDB 连接器http://docs.mongodb.org/ecosystem/tools/hadoop/?_ga=1.111209414.370990604.1441913822
  2. 实现您自己的 OutputFormathttps://hadoop.apache.org/docs/r2.7.0/api/org/apache/hadoop/mapred/OutputFormat.html(改为使用 FileOutputFormat)。
  3. 在 reducer 中执行 mongodb 查询,而不是在 MapREduce 上下文中写入(不好,根据驱动程序中指定的 OutputFormat,您可能会以 HDFS 中的空输出文件结束)

在我看来,选项 1 是最好的选择,但我没有使用 MongoDB 连接器来说明它是否足够稳定和功能强大。选项 2 要求您真正了解 hadoop underhood 是如何工作的,以避免出现大量打开的连接以及事务和 hadoop 任务重试问题。

【讨论】:

  • Rojo 你知道如何将输出数据存储在对象中而不是存储到文件系统中
  • 你能详细说明你的问题吗?如果是与原始问题不同的问题,您能否创建一个新问题,让有相同问题的其他人更容易找到它?如果我知道答案,我会尽力帮助你
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-01-26
  • 2018-05-06
  • 1970-01-01
  • 2013-06-03
  • 1970-01-01
相关资源
最近更新 更多