diff --git a/hadoop-hdfs-project/hadoop-hdfs-client/src/main/java/org/apache/hadoop/hdfs/StripeReader.java b/hadoop-hdfs-project/hadoop-hdfs-client/src/main/java/org/apache/hadoop/hdfs/StripeReader.java index f2d6732a4596de..bc39bace795881 100644 --- a/hadoop-hdfs-project/hadoop-hdfs-client/src/main/java/org/apache/hadoop/hdfs/StripeReader.java +++ b/hadoop-hdfs-project/hadoop-hdfs-client/src/main/java/org/apache/hadoop/hdfs/StripeReader.java @@ -253,6 +253,9 @@ private int readToBuffer(BlockReader blockReader, strategy.getReadBuffer().clear(); // we want to remember which block replicas we have tried corruptedBlocks.addCorruptedBlock(currentBlock, currentNode); + if (blockReader != null) { + blockReader.close(); + } throw ce; } catch (IOException e) { DFSClient.LOG.warn("Exception while reading from " @@ -260,6 +263,9 @@ private int readToBuffer(BlockReader blockReader, + currentNode, e); //Clear buffer to make next decode success strategy.getReadBuffer().clear(); + if (blockReader != null) { + blockReader.close(); + } throw e; } } @@ -329,21 +335,26 @@ boolean readChunk(final LocatedBlock block, int chunkIndex) * read the whole stripe. do decoding if necessary */ void readStripe() throws IOException { - for (int i = 0; i < dataBlkNum; i++) { - if (alignedStripe.chunks[i] != null && - alignedStripe.chunks[i].state != StripingChunk.ALLZERO) { - if (!readChunk(targetBlocks[i], i)) { - alignedStripe.missingChunksNum++; + try { + for (int i = 0; i < dataBlkNum; i++) { + if (alignedStripe.chunks[i] != null && + alignedStripe.chunks[i].state != StripingChunk.ALLZERO) { + if (!readChunk(targetBlocks[i], i)) { + alignedStripe.missingChunksNum++; + } } } - } - // There are missing block locations at this stage. Thus we need to read - // the full stripe and one more parity block. - if (alignedStripe.missingChunksNum > 0) { - checkMissingBlocks(); - readDataForDecoding(); - // read parity chunks - readParityChunks(alignedStripe.missingChunksNum); + // There are missing block locations at this stage. Thus we need to read + // the full stripe and one more parity block. + if (alignedStripe.missingChunksNum > 0) { + checkMissingBlocks(); + readDataForDecoding(); + // read parity chunks + readParityChunks(alignedStripe.missingChunksNum); + } + } catch (IOException e) { + dfsStripedInputStream.close(); + throw e; } // TODO: for a full stripe we can start reading (dataBlkNum + 1) chunks @@ -385,7 +396,8 @@ void readStripe() throws IOException { } } catch (InterruptedException ie) { String err = "Read request interrupted"; - DFSClient.LOG.error(err); + DFSClient.LOG.error(err, ie); + dfsStripedInputStream.close(); clearFutures(); // Don't decode if read interrupted throw new InterruptedIOException(err); diff --git a/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/LineRecordReader.java b/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/LineRecordReader.java index ab63c199f2f35b..4fd9a5fda0a551 100644 --- a/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/LineRecordReader.java +++ b/hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/LineRecordReader.java @@ -153,7 +153,12 @@ public LineRecordReader(Configuration job, FileSplit split, // because we always (except the last split) read one extra line in // next() method. if (start != 0) { - start += in.readLine(new Text(), 0, maxBytesToConsume(start)); + try { + start += in.readLine(new Text(), 0, maxBytesToConsume(start)); + } catch (Exception e) { + close(); + throw e; + } } this.pos = start; }