如何在Eclipse中无需输入文件,基于对象列表使用Hadoop MapReduce
嘿,这个需求完全可以实现!核心就是自定义一套不依赖文件的InputFormat体系,绕开Hadoop默认对输入文件的依赖,直接用你内存中的对象列表生成键值对。我帮你把完整的方案和代码补全,直接在Eclipse里就能跑:
1. 自定义InputSplit
这个类只是一个标记,因为我们的数据源是内存列表,不需要拆分文件,所以实现最基础的序列化逻辑就行:
import org.apache.hadoop.io.Writable; import org.apache.hadoop.mapreduce.InputSplit; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; public class ObjectListInputSplit extends InputSplit implements Writable { private int splitId; // 必须要有空构造函数,Hadoop序列化需要 public ObjectListInputSplit() {} public ObjectListInputSplit(int splitId) { this.splitId = splitId; } @Override public long getLength() throws IOException, InterruptedException { // 非文件数据源,返回0即可 return 0; } @Override public String[] getLocations() throws IOException, InterruptedException { // 不需要指定数据节点,返回空数组 return new String[0]; } @Override public void write(DataOutput out) throws IOException { out.writeInt(splitId); } @Override public void readFields(DataInput in) throws IOException { splitId = in.readInt(); } }
2. 自定义RecordReader(核心)
这个类负责从你的对象列表中读取数据,生成MapReduce需要的键值对。测试阶段我们用静态变量传递对象列表,简单直接:
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.RecordReader; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.hadoop.mapreduce.InputSplit; import java.io.IOException; import java.util.Iterator; import java.util.List; public class ObjectListRecordReader extends RecordReader<Text, Text> { // 替换成你自己的对象类型,比如你的自定义Object类 private static List<String> objectList; private Iterator<String> iterator; private Text currentKey; private Text currentValue; // 静态方法用来传入你的数据源列表 public static void setObjectList(List<String> list) { objectList = list; } @Override public void initialize(InputSplit split, TaskAttemptContext context) throws IOException, InterruptedException { // 初始化列表迭代器 iterator = objectList.iterator(); } @Override public boolean nextKeyValue() throws IOException, InterruptedException { if (iterator.hasNext()) { String obj = iterator.next(); // 这里替换成你的业务逻辑:从对象生成Key和Value currentKey = new Text("custom_key_" + obj); currentValue = new Text(obj); return true; } currentKey = null; currentValue = null; return false; } @Override public Text getCurrentKey() throws IOException, InterruptedException { return currentKey; } @Override public Text getCurrentValue() throws IOException, InterruptedException { return currentValue; } @Override public float getProgress() throws IOException, InterruptedException { // 简单计算处理进度,按需调整 return 1.0f - (iterator.hasNext() ? 1.0f / objectList.size() : 0); } @Override public void close() throws IOException { // 无资源需要关闭,空实现即可 } }
3. 自定义InputFormat
负责创建我们的自定义Split和RecordReader,这里只生成一个Split(如果需要并行处理,可以拆分多个Split):
import org.apache.hadoop.mapreduce.InputFormat; import org.apache.hadoop.mapreduce.JobContext; import org.apache.hadoop.mapreduce.RecordReader; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.hadoop.mapreduce.InputSplit; import java.util.ArrayList; import java.util.List; public class ObjectListInputFormat extends InputFormat<Text, Text> { @Override public List<InputSplit> getSplits(JobContext context) throws IOException, InterruptedException { List<InputSplit> splits = new ArrayList<>(); // 生成一个Split,如需并行可添加多个 splits.add(new ObjectListInputSplit(0)); return splits; } @Override public RecordReader<Text, Text> createRecordReader(InputSplit split, TaskAttemptContext context) throws IOException, InterruptedException { return new ObjectListRecordReader(); } }
4. 主类中使用
在Job配置里指定我们的自定义InputFormat,然后传入你的对象列表即可:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.util.ArrayList; import java.util.List; public class ObjectListMapReduce { public static class MyMapper extends Mapper<Text, Text, Text, Text> { @Override protected void map(Text key, Text value, Context context) throws IOException, InterruptedException { // 你的Map逻辑,和普通MapReduce完全一致 context.write(key, value); } } public static class MyReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 你的Reduce逻辑 for (Text val : values) { context.write(key, val); } } } public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "ObjectListMapReduceDemo"); job.setJarByClass(ObjectListMapReduce.class); // 关键:设置自定义InputFormat job.setInputFormatClass(ObjectListInputFormat.class); // 准备你的对象列表,替换成你自己的对象集合 List<String> myObjectList = new ArrayList<>(); myObjectList.add("user_001"); myObjectList.add("user_002"); myObjectList.add("user_003"); // 把列表传给RecordReader ObjectListRecordReader.setObjectList(myObjectList); job.setMapperClass(MyMapper.class); job.setReducerClass(MyReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // 设置输出路径(MapReduce仍需输出到文件系统,如需自定义输出可再写OutputFormat) FileOutputFormat.setOutputPath(job, new Path(args[0])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
注意事项
- 上面用静态变量传递列表的方式,在Eclipse本地运行完全没问题;如果要部署到集群,需要把对象序列化成字节数组存入
Configuration,或者用分布式缓存传递。 - 如果你的对象不是
String,要确保对象是可序列化的,或者在RecordReader里自己处理对象的读取和键值转换逻辑。 - 不需要设置任何输入路径,因为我们的
InputFormat完全不依赖文件。
内容的提问来源于stack exchange,提问作者Henry Daniel Saenz
相关产品推荐
相关产品推荐

