You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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租约机制冲突:

  1. 当多个Reducer任务并行运行时,所有Reducer都会在cleanup阶段尝试删除并重建同一个centroids.seq文件。第一个Reducer删除文件后创建新文件时,其他Reducer可能仍持有旧文件的租约,或者在尝试写入时发现文件已被修改,导致NameNode验证租约失败。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 06:17:03