diff --git a/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordReader.java b/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordReader.java index fc90e5ccd4..73f978683a 100644 --- a/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordReader.java +++ b/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordReader.java @@ -106,7 +106,9 @@ public class ColumnFamilyRecordReader extends RecordReader 1.0F ? 1.0F : progress; @@ -423,6 +425,7 @@ public class ColumnFamilyRecordReader extends RecordReader>> wideColumns; private ByteBuffer lastColumn = ByteBufferUtil.EMPTY_BYTE_BUFFER; + private ByteBuffer lastCountedKey = ByteBufferUtil.EMPTY_BYTE_BUFFER; private void maybeInit() { @@ -476,12 +479,28 @@ public class ColumnFamilyRecordReader extends RecordReader> next = wideColumns.next(); lastColumn = next.right.values().iterator().next().name(); + + maybeCountRow(next); return next; } + + /** + * Increases the row counter only if we really moved to the next row. + * @param next just fetched row slice + */ + private void maybeCountRow(Pair> next) + { + ByteBuffer currentKey = next.left; + if (!currentKey.equals(lastCountedKey)) + { + totalRead++; + lastCountedKey = currentKey; + } + } + private class WideColumnIterator extends AbstractIterator>> { private final Iterator rows;