ラベル hadoop の投稿を表示しています。 すべての投稿を表示
ラベル hadoop の投稿を表示しています。 すべての投稿を表示

2011年3月16日水曜日

Hadoop のHDFS上のファイルを読む

サイドファイルにObjectOutputStreamで書いたものはObjectInputStreamで読めばいい。 dirname 以下のfilterにひっかかるファイルからオブジェクトを読み出すのはこうする。
  Configuration conf = new Configuration();
  FileSystem fs = FileSystem.get(URI.create("hdfs://localhost:9000/"), conf);
  for (FileStatus status: fs.listStatus(new Path(dirname), filter)) {
   InputStream is = fs.open(status.getPath());
   ObjectInputStream ois = new ObjectInputStream(is); 
   Object o = ois.readObject();
   System.out.println(o);
  }
filter はこんな風に取得。filePrefix に指定した文字を含むファイルを選ぶ。この実装だとprefix になってないけど。
 static PathFilter getPathFilter(final String filePrefix){
  PathFilter filter = new PathFilter() {
   public boolean accept(Path path) {
    return path.getName().contains(filePrefix);
   }
  };
  return filter;
 }

普通に書き出したファイルの場合

サイドデータじゃなくて、普通にcontext.writeした場合は読み方が変わってくる。SequenceFile.Readerで読む。
  Configuration conf = new Configuration();

  FileSystem fs = FileSystem.get(URI.create("hdfs://localhost:9000/"), conf);
  for (FileStatus status: fs.listStatus(new Path(dirname), filter)) {
   SequenceFile.Reader reader = new SequenceFile.Reader(fs, status.getPath(), conf);
   Text key = new Text();
                        Text value = new Text();
   while (reader.next(key, value){
                               ....
   }
  }
この読み方は、ValueがWritableの時にしか使えない。Serializable の場合はつぎのようにする。 nextでkeyだけ読んで、getCurrentValueでvalueを読む。このときにconfにJavaSerializationを追加しておかないと、エラーになるので注意。
  Configuration conf = new Configuration();
  conf.set("io.serializations", 
    JavaSerialization.class.getName() + "," +
    WritableSerialization.class.getName());

  FileSystem fs = FileSystem.get(URI.create("hdfs://localhost:9000/"), conf);
  for (FileStatus status: fs.listStatus(new Path(dirname), filter)) {
   SequenceFile.Reader reader = new SequenceFile.Reader(fs, status.getPath(), conf);
   Text key = new Text();
   while (reader.next(key)){
    System.out.println(key.toString());
    ChlacTest.Result res = new ChlacTest.Result();
    res = (Result) reader.getCurrentValue(res);
    System.out.println(String.format("%d %d %f", res.rx, res.time_frame, res.alpha));
   }
  }

2011年3月11日金曜日

hadoop side effect file

hadoopではファイルを正式なアウトプットの他に作ることができる。 getWorkOutputPath を使う。ディレクトリは正式なアウトプットの出力先と同じなので注意。
 public static class R1 extends
 Reducer {
  public void reduce(Text key, Iterable values,
  Context context) throws IOException, InterruptedException {
   context.write(key, values.iterator().next());
   Path path = FileOutputFormat.getWorkOutputPath(context); 
   Path fpath = new Path(path, "sidefile" + key.toString());  
    OutputStream os = fpath.getFileSystem(context.getConfiguration()).create(fpath);
   PrintWriter pw = new PrintWriter(new OutputStreamWriter(os));
   pw.println("hello");
   pw.flush();
   pw.close();
   os.close();
  }
 }

2011年3月9日水曜日

Hadoop job 管理

ps 相当
> hadoop job -list
kill
> hadoop job -kill xxxxxxxx

2011年2月4日金曜日

hadoop でjava のserializableを受け渡すには

何らかの方法でシリアライズするわけだが、hadoop はデフォルトではJavaのシリアライズではなく 独自のシリアライザを使うようになっている。これはJavaのシリアライズが重いため、だそうだ。

しかしもちろんシリアライザの実装を変更することができ、Javaのシリアライザを使うように 指定することもできる。それには、こうする。

Configuration conf = new Configuration();
conf.set("io.serializations", 
  org.apache.hadoop.io.serializer.JavaSerialization.class.getName() + "," +
  org.apache.hadoop.io.serializer.WritableSerialization.class.getName());
要するにSerialization クラスをカンマで区切って、 io.serializations にセットするのだけど、JavaSerialization だけだと、Textとかが デコードできなくなっちゃうので、デフォルトのWritableSerialization も書いておくこと。