Hadoop执行Kmeans程序时文件存在却报FileNotFoundException求助
解决Hadoop Kmeans程序中的FileNotFoundException问题
执行Hadoop Kmeans程序命令:
hadoop jar kmeans-1.0-SNAPSHOT.jar it.kurapika.Kmeans dataset.txt output
出现如下错误,尽管通过Hadoop门户确认/user/hadoop/centroids.seq文件存在:
Error: java.io.FileNotFoundException: File does not exist: /user/hadoop/centroids.seq (inode 17689) [Lease. Holder: DFSClient_attempt_1685688288161_0036_r_000002_0_2097879656_1, pending creates: 1] at org.apache.hadoop.hdfs.server.namenode.FSNamesystem.checkLease(FSNamesystem.java:2840) at org.apache.hadoop.hdfs.server.namenode.FSDirWriteFileOp.analyzeFileState(FSDirWriteFileOp.java:599) at org.apache.hadoop.hdfs.server.namenode.FSDirWriteFileOp.validateAddBlock(FSDirWriteFileOp.java:171) at org.apache.hadoop.hdfs.server.namenode.FSNamesystem.getAdditionalBlock(FSNamesystem.java:2719) at org.apache.hadoop.hdfs.server.namenode.NameNodeRpcServer.addBlock(NameNodeRpcServer.java:892) at org.apache.hadoop.hdfs.protocolPB.ClientNamenodeProtocolServerSideTranslatorPB.addBlock(ClientNamenodeProtocolServerSideTranslatorPB.java:568) at org.apache.hadoop.hdfs.protocol.proto.ClientNamenodeProtocolProtos$ClientNamenodeProtocol$2.callBlockingMethod(ClientNamenodeProtocolProtos.java) at org.apache.hadoop.ipc.ProtobufRpcEngine$Server$ProtoBufRpcInvoker.call(ProtobufRpcEngine.java:527) at org.apache.hadoop.ipc.RPC$Server.call(RPC.java:1036) at org.apache.hadoop.ipc.Server$RpcCall.run(Server.java:1000) at org.apache.hadoop.ipc.Server$RpcCall.run(Server.java:928) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1729) at org.apache.hadoop.ipc.Server$Handler.run(Server.java:2916) at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method) at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62) at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45) at java.lang.reflect.Constructor.newInstance(Constructor.java:423) at org.apache.hadoop.ipc.RemoteException.instantiateException(RemoteException.java:121) at org.apache.hadoop.ipc.RemoteException.unwrapRemoteException(RemoteException.java:88) at org.apache.hadoop.hdfs.DFSOutputStream.addBlock(DFSOutputStream.java:1084) at org.apache.hadoop.hdfs.DataStreamer.locateFollowingBlock(DataStreamer.java:1866) at org.apache.hadoop.hdfs.DataStreamer.nextBlockOutputStream(DataStreamer.java:1668) at org.apache.hadoop.hdfs.DataStreamer.run(DataStreamer.java:716)
程序逻辑:Driver类Kmeans.java创建centroids.seq文件,Reducer的cleanup方法删除该文件并重新写入新质心数据。尝试执行hadoop fs -chmod +wx /user/hadoop/centroids.seq修改权限后问题未解决。
相关代码
Kmeans.java
public class Kmeans { public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException { Configuration conf = new Configuration(); conf.addResource(new Path("configuration.xml")); String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs(); // set parameters final String INPUT = otherArgs[0]; final String OUTPUT = otherArgs[1] + "/temp"; final int DATASET_SIZE = conf.getInt("dataset.size", 100); final String CENTROIDS_PATH = conf.get("centroids.path", "centroids.seq"); final int K = conf.getInt("k", 3); final int MAX_ITERATIONS = conf.getInt("iterations", 20); Point[] newCentroids = new Point[K]; // generate initial centroids newCentroids = Utility.generateCentroids(conf, INPUT, K, DATASET_SIZE); Utility.writeCentroids(conf, new Path(CENTROIDS_PATH), newCentroids); ... // write final centroids in output file Utility.writeOutput(conf, new Path(CENTROIDS_PATH), new Path(otherArgs[1])); System.exit(0); } }
Utility.java
public class Utility { public static Point[] generateCentroids(Configuration conf, String pathString, int k, int dataSetSize) throws IOException { Point[] points = new Point[k]; //Create a sorted list of positions without duplicates //Positions are the line index of the random selected centroids List<Integer> positions = new ArrayList<Integer>(); Random random = new Random(); int pos; while(positions.size() < k) { pos = random.nextInt(dataSetSize); if(!positions.contains(pos)) { positions.add(pos); } } Collections.sort(positions); //File reading utils Path dataPath = new Path(pathString); FileSystem hdfs = FileSystem.get(conf); FSDataInputStream in = hdfs.open(dataPath); BufferedReader br = new BufferedReader(new InputStreamReader(in)); //Get centroids from the file int row = 0; int i = 0; int position; while(i < positions.size()) { position = positions.get(i); String point = br.readLine(); if(row == position) { points[i] = new Point(); points[i].parse(point); i++; } row++; } br.close(); return points; } public static void writeCentroids(Configuration conf, Path center, Point[] points) throws IOException { try (SequenceFile.Writer centerWriter = SequenceFile.createWriter(conf, SequenceFile.Writer.file(center) , SequenceFile.Writer.keyClass(Point.class), SequenceFile.Writer.valueClass(IntWritable.class))) { final IntWritable value = new IntWritable(0); for (Point point : points) { centerWriter.append(point, value); } } } public static void writeOutput(Configuration conf, Path centroidsPath, Path outpath) throws IOException { List<Centroid> centroids = new ArrayList<>(); // list of centroids try (SequenceFile.Reader reader = new SequenceFile.Reader(conf, SequenceFile.Reader.file(centroidsPath))) { Centroid key = new Centroid(); while (reader.next(key)) { // iterate over records Centroid center = new Centroid(key); // create new centroid centroids.add(center); // add new Centroid to list } FileSystem hdfs = FileSystem.get(conf); FSDataOutputStream dos = hdfs.create(outpath, true); BufferedWriter br = new BufferedWriter(new OutputStreamWriter(dos)); // write the result in output file for(int i = 0; i < centroids.size(); i++) { br.write(centroids.get(i).toString()); br.newLine(); } br.close(); hdfs.close(); } } }
Reducer类
public class KmeansReducer extends Reducer<Centroid, Point, Centroid, NullWritable>{ public static enum Counter { // Global counter: it gets incremented every time new centroids are more than epsilon distant from previous centroids CONVERGED } private final List<Centroid> centers = new ArrayList<>(); // list containing new centroids private Double epsilon = 0.; // convergence parameter @Override protected void setup(Context context) { Configuration conf = context.getConfiguration(); epsilon = conf.getDouble("epsilon", 0.0001); // initialize convergence parameter with value in configuration file } // for each cluster calculate new centroids @Override protected void reduce(Centroid key, Iterable<Point> partialSums, Context context) throws IOException, InterruptedException { Centroid newKey = new Centroid(); // new centroid for (Point point : partialSums) { // summation of partial sums newKey.getPoint().sum(point); } newKey.getPoint().compress(); // divide for number of points in cluster newKey.setIndex(key); // assign old centroid's index to new centroid centers.add(newKey); // add new centroid to new centroids list context.write(newKey, NullWritable.get()); // write output record (key: centroid, value: null) // calculate distance between new centroid and old centroid double distance = key.getPoint().getDistance(newKey.getPoint()); if (distance > epsilon) { // if distance is greater than epsilon context.getCounter(Counter.CONVERGED).increment(1); // increment global counter } } // write new centroids in sequence file @Override protected void cleanup(Context context) throws IOException, InterruptedException { Configuration conf = context.getConfiguration(); Path outPath = new Path(conf.get("centroids.path", "centroids.seq")); // get path of centroids sequence file FileSystem fs = FileSystem.get(conf); fs.delete(outPath, true); // if path exists delete it try (SequenceFile.Writer out = SequenceFile.createWriter(conf, SequenceFile.Writer.file(outPath), SequenceFile.Writer.keyClass(Centroid.class), SequenceFile.Writer.valueClass(IntWritable.class))) { final IntWritable value = new IntWritable(0); for (Centroid center : centers) { out.append(center, value); // write new centroids in sequence file } } } }
问题原因
这个错误的核心是HDFS租约机制冲突:
- 当多个Reducer任务并行运行时,所有Reducer都会在
cleanup阶段尝试删除并重建同一个centroids.seq文件。第一个Reducer删除文件后创建新文件时,其他Reducer可能仍持有旧文件的租约,或者在尝试写入时发现文件已被修改,导致NameNode验证租约失败。 - HDFS的文件租约是为了防止多个客户端同时修改同一个文件,当一个客户端持有租约时,其他客户端无法对文件进行写操作。Reducer直接操作全局文件的方式违反了HDFS的并发安全模型。
解决方案
1. 重构Reducer逻辑,避免直接操作全局文件
去掉Reducer的cleanup方法中对centroids.seq的操作,改为将新质心输出到MapReduce的Context中,由框架写入临时输出目录:
// 修改Reducer的cleanup方法,删除原文件操作代码 @Override protected void cleanup(Context context) throws IOException, InterruptedException { // 直接将所有新质心输出到Context for (Centroid center : centers) { context.write(center, NullWritable.get()); } }
2. 由Driver统一处理质心更新
在Driver的迭代逻辑中,每次MapReduce任务完成后,从临时输出目录读取所有Reducer生成的新质心,合并计算出全局质心,再写入centroids.seq文件:
// 在Kmeans.java的迭代循环中添加以下逻辑 for (int iter = 0; iter < MAX_ITERATIONS; iter++) { // 启动MapReduce任务 Job job = Job.getInstance(conf, "Kmeans Iteration " + iter); // ... 设置Job的Mapper、Reducer、输入输出等参数 ... job.waitForCompletion(true); // 检查是否收敛 long convergedCount = job.getCounters().findCounter(KmeansReducer.Counter.CONVERGED).getValue(); if (convergedCount == 0) { break; } // 从临时输出目录读取所有新质心 Path tempOutput = new Path(OUTPUT + "_" + iter); List<Centroid> partialCentroids = Utility.readPartialCentroids(conf, tempOutput); // 合并计算全局质心(根据每个质心的index分组求和) Map<Integer, Centroid> globalCentroidsMap = new HashMap<>(); for (Centroid c : partialCentroids) { int index = c.getIndex(); if (globalCentroidsMap.containsKey(index)) { globalCentroidsMap.get(index).getPoint().sum(c.getPoint()); } else { globalCentroidsMap.put(index, new Centroid(c)); } } // 计算平均得到最终新质心 Point[] newCentroids = new Point[K]; for (Map.Entry<Integer, Centroid> entry : globalCentroidsMap.entrySet()) { entry.getValue().getPoint().compress(); newCentroids[entry.getKey()] = entry.getValue().getPoint(); } // 写入新质心到centroids.seq Utility.writeCentroids(conf, new Path(CENTROIDS_PATH), newCentroids); // 删除临时输出目录 FileSystem.get(conf).delete(tempOutput, true); }
3. 确保文件操作原子性
如果必须直接操作文件,使用HDFS的原子重命名操作:先写入临时文件,再将临时文件重命名为目标路径,避免中间状态被其他任务读取:
// 修改Utility.writeCentroids方法 public static void writeCentroids(Configuration conf, Path center, Point[] points) throws IOException { // 创建临时文件 Path tempPath = new Path(center.toString() + ".tmp"); try (SequenceFile.Writer centerWriter = SequenceFile.createWriter(conf, SequenceFile.Writer.file(tempPath) , SequenceFile.Writer.keyClass(Point.class), SequenceFile.Writer.valueClass(IntWritable.class))) { final IntWritable value = new IntWritable(0); for (Point point : points) { centerWriter.append(point, value); } } // 原子重命名为目标文件 FileSystem fs = FileSystem.get(conf); if (fs.exists(center)) { fs.delete(center, true); } fs.rename(tempPath, center); }
4. 检查租约释放
确保所有文件流都被正确关闭,使用try-with-resources语法(代码中已使用,保持即可),避免租约被长期占用。
内容的提问来源于stack exchange,提问作者Sissi
相关产品推荐
相关产品推荐

