From d8b1fc3952caa22be53e20d0e145aa3137ffc58e Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 16 Jun 2011 04:24:33 +0000 Subject: [PATCH] replace CollatingIterator, ReducingIterator with MergeIterator patch by stuhood; reviewed by jbellis for CASSANDRA-2062 git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@1136287 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 + .../cassandra/db/ColumnFamilyStore.java | 12 +- .../org/apache/cassandra/db/RowIterator.java | 72 --------- .../cassandra/db/RowIteratorFactory.java | 67 ++++---- .../db/columniterator/IColumnIterator.java | 3 +- .../db/compaction/CompactionIterator.java | 148 ++++++++++-------- .../db/compaction/CompactionManager.java | 13 +- .../db/compaction/LazilyCompactedRow.java | 31 ++-- .../cassandra/db/filter/QueryFilter.java | 19 +-- .../cassandra/io/sstable/KeyIterator.java | 4 +- .../io/sstable/ReducingKeyIterator.java | 35 ++--- .../cassandra/io/sstable/SSTableScanner.java | 5 +- .../service/RangeSliceResponseResolver.java | 42 ++--- .../cassandra/service/RowRepairResolver.java | 9 +- .../apache/cassandra/utils/FBUtilities.java | 51 +++--- .../cassandra/utils/ReducingIterator.java | 87 ---------- .../cassandra/utils/MergeIteratorTest.java | 107 +++++++++++++ 17 files changed, 321 insertions(+), 386 deletions(-) delete mode 100644 src/java/org/apache/cassandra/db/RowIterator.java delete mode 100644 src/java/org/apache/cassandra/utils/ReducingIterator.java create mode 100644 test/unit/org/apache/cassandra/utils/MergeIteratorTest.java diff --git a/CHANGES.txt b/CHANGES.txt index e4bd3358c7..6e1ee1513d 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -3,6 +3,8 @@ * add commitlog_total_space_in_mb to prevent fragmented logs (CASSANDRA-2427) * removed commitlog_rotation_threshold_in_mb configuration (CASSANDRA-2771) * make AbstractBounds.normalize de-overlapp overlapping ranges (CASSANDRA-2641) + * replace CollatingIterator, ReducingIterator with MergeIterator + (CASSANDRA-2062) 0.8.1 diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index b536dbbd4e..9a82ee770e 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -32,7 +32,6 @@ import javax.management.MBeanServer; import javax.management.ObjectName; import com.google.common.collect.Iterables; -import org.apache.commons.collections.IteratorUtils; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -1188,7 +1187,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean IColumnIterator ci = filter.getMemtableColumnIterator(cached, null, getComparator()); ColumnFamily cf = ci.getColumnFamily().cloneMeShallow(); - filter.collectCollatedColumns(cf, ci, gcBefore); + filter.collateColumns(cf, Collections.singletonList(ci), getComparator(), gcBefore); // TODO this is necessary because when we collate supercolumns together, we don't check // their subcolumns for relevance, so we need to do a second prune post facto here. return cf.isSuper() ? removeDeleted(cf, gcBefore) : removeDeletedCF(cf, gcBefore); @@ -1244,10 +1243,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean if (iterators.size() == 0) return null; - Comparator comparator = filter.filter.getColumnComparator(getComparator()); - Iterator collated = IteratorUtils.collatedIterator(comparator, iterators); - - filter.collectCollatedColumns(returnCF, collated, gcBefore); + filter.collateColumns(returnCF, iterators, getComparator(), gcBefore); // Caller is responsible for final removeDeletedCF. This is important for cacheRow to work correctly: return returnCF; @@ -1298,7 +1294,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean // It is fine to aliases the View.sstables since it's an unmodifiable collection Collection sstables = currentView.sstables; - RowIterator iterator = RowIteratorFactory.getIterator(memtables, sstables, startWith, stopAt, filter, getComparator(), this); + CloseableIterator iterator = RowIteratorFactory.getIterator(memtables, sstables, startWith, stopAt, filter, getComparator(), this); List rows = new ArrayList(); try @@ -1486,7 +1482,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean ColumnFamily expandedData = data; data = expandedData.cloneMeShallow(); IColumnIterator iter = dataFilter.getMemtableColumnIterator(expandedData, dk, getComparator()); - new QueryFilter(dk, path, dataFilter).collectCollatedColumns(data, iter, gcBefore()); + new QueryFilter(dk, path, dataFilter).collateColumns(data, Collections.singletonList(iter), getComparator(), gcBefore()); } rows.add(new Row(dk, data)); diff --git a/src/java/org/apache/cassandra/db/RowIterator.java b/src/java/org/apache/cassandra/db/RowIterator.java deleted file mode 100644 index 410c8ec4b5..0000000000 --- a/src/java/org/apache/cassandra/db/RowIterator.java +++ /dev/null @@ -1,72 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.cassandra.db; - -import java.io.Closeable; -import java.io.IOException; -import java.util.Iterator; -import java.util.List; - -import org.apache.cassandra.db.columniterator.IColumnIterator; -import org.apache.cassandra.utils.ReducingIterator; - -/** - * Row iterator that allows us to close the underlying iterators. - */ -public class RowIterator implements Closeable, Iterator -{ - private ReducingIterator reduced; - private List> iterators; - - /** - * @param reduced Reducing iterator that takes multiple iterators and provides us with - * one row at the time. - * @param iterators The underlying iterators that we will close when done. - */ - public RowIterator(ReducingIterator reduced, List> iterators) - { - this.reduced = reduced; - this.iterators = iterators; - } - - public boolean hasNext() - { - return reduced.hasNext(); - } - - public Row next() - { - return reduced.next(); - } - - public void remove() - { - reduced.remove(); - } - - public void close() throws IOException - { - for (Iterator iter : iterators) - { - if (iter instanceof Closeable) - { - ((Closeable)iter).close(); - } - } - } -} diff --git a/src/java/org/apache/cassandra/db/RowIteratorFactory.java b/src/java/org/apache/cassandra/db/RowIteratorFactory.java index 5764b43693..c603f7e431 100644 --- a/src/java/org/apache/cassandra/db/RowIteratorFactory.java +++ b/src/java/org/apache/cassandra/db/RowIteratorFactory.java @@ -25,15 +25,16 @@ import java.util.Map.Entry; import com.google.common.base.Function; import com.google.common.base.Predicate; +import com.google.common.collect.AbstractIterator; import com.google.common.collect.Iterators; -import org.apache.commons.collections.IteratorUtils; import org.apache.cassandra.db.columniterator.IColumnIterator; import org.apache.cassandra.db.filter.QueryFilter; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.io.sstable.SSTableScanner; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.MergeIterator; public class RowIteratorFactory { @@ -60,7 +61,7 @@ public class RowIteratorFactory * @param comparator * @return A row iterator following all the given restrictions */ - public static RowIterator getIterator(final Collection memtables, + public static CloseableIterator getIterator(final Collection memtables, final Collection sstables, final DecoratedKey startWith, final DecoratedKey stopAt, @@ -70,7 +71,7 @@ public class RowIteratorFactory ) { // fetch data from current memtable, historical memtables, and SSTables in the correct order. - final List> iterators = new ArrayList>(); + final List> iterators = new ArrayList>(); // we iterate through memtables with a priority queue to avoid more sorting than necessary. // this predicate throws out the rows before the start of our range. Predicate p = new Predicate() @@ -85,8 +86,7 @@ public class RowIteratorFactory // memtables for (Memtable memtable : memtables) { - iterators.add(Iterators.filter(Iterators.transform(memtable.getEntryIterator(startWith), - new ConvertToColumnIterator(filter, comparator)), p)); + iterators.add(new ConvertToColumnIterator(filter, comparator, p, memtable.getEntryIterator(startWith))); } // sstables @@ -98,10 +98,9 @@ public class RowIteratorFactory iterators.add(scanner); } - Iterator collated = IteratorUtils.collatedIterator(COMPARE_BY_KEY, iterators); - + final Memtable firstMemtable = memtables.iterator().next(); // reduce rows from all sources into a single row - ReducingIterator reduced = new ReducingIterator(collated) + return MergeIterator.get(iterators, COMPARE_BY_KEY, new MergeIterator.Reducer() { private final int gcBefore = (int) (System.currentTimeMillis() / 1000) - cfs.metadata.getGcGraceSeconds(); private final List colIters = new ArrayList(); @@ -121,57 +120,61 @@ public class RowIteratorFactory this.returnCF.delete(current.getColumnFamily()); } - @Override - protected boolean isEqual(IColumnIterator o1, IColumnIterator o2) - { - return COMPARE_BY_KEY.compare(o1, o2) == 0; - } - protected Row getReduced() { - Comparator colComparator = filter.filter.getColumnComparator(comparator); - Iterator colCollated = IteratorUtils.collatedIterator(colComparator, colIters); // First check if this row is in the rowCache. If it is we can skip the rest ColumnFamily cached = cfs.getRawCachedRow(key); - if (cached != null) + if (cached == null) + // not cached: collate + filter.collateColumns(returnCF, colIters, comparator, gcBefore); + else { QueryFilter keyFilter = new QueryFilter(key, filter.path, filter.filter); returnCF = cfs.filterColumnFamily(cached, keyFilter, gcBefore); } - else if (colCollated.hasNext()) - { - filter.collectCollatedColumns(returnCF, colCollated, gcBefore); - } Row rv = new Row(key, returnCF); colIters.clear(); key = null; return rv; } - }; - - return new RowIterator(reduced, iterators); + }); } /** * Get a ColumnIterator for a specific key in the memtable. */ - private static class ConvertToColumnIterator implements Function, IColumnIterator> + private static class ConvertToColumnIterator extends AbstractIterator implements CloseableIterator { - private QueryFilter filter; - private AbstractType comparator; + private final QueryFilter filter; + private final AbstractType comparator; + private final Predicate pred; + private final Iterator> iter; - public ConvertToColumnIterator(QueryFilter filter, AbstractType comparator) + public ConvertToColumnIterator(QueryFilter filter, AbstractType comparator, Predicate pred, Iterator> iter) { this.filter = filter; this.comparator = comparator; + this.pred = pred; + this.iter = iter; } - public IColumnIterator apply(final Entry entry) + public IColumnIterator computeNext() { - return filter.getMemtableColumnIterator(entry.getValue(), entry.getKey(), comparator); + while (iter.hasNext()) + { + Map.Entry entry = iter.next(); + IColumnIterator ici = filter.getMemtableColumnIterator(entry.getValue(), entry.getKey(), comparator); + if (pred.apply(ici)) + return ici; + } + return endOfData(); + } + + public void close() + { + // pass } } - } diff --git a/src/java/org/apache/cassandra/db/columniterator/IColumnIterator.java b/src/java/org/apache/cassandra/db/columniterator/IColumnIterator.java index a0d93ca65b..1adaa6ba54 100644 --- a/src/java/org/apache/cassandra/db/columniterator/IColumnIterator.java +++ b/src/java/org/apache/cassandra/db/columniterator/IColumnIterator.java @@ -27,8 +27,9 @@ import java.util.Iterator; import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.IColumn; +import org.apache.cassandra.utils.CloseableIterator; -public interface IColumnIterator extends Iterator +public interface IColumnIterator extends CloseableIterator { /** * @return An empty CF holding metadata for the row being iterated. diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index 4b1013b4bf..ff6075f17c 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -21,33 +21,35 @@ package org.apache.cassandra.db.compaction; */ -import java.io.Closeable; import java.io.IOException; import java.util.ArrayList; +import java.util.Comparator; import java.util.Iterator; import java.util.List; -import org.apache.cassandra.service.StorageService; -import org.apache.commons.collections.iterators.CollatingIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.google.common.collect.AbstractIterator; + import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.columniterator.IColumnIterator; import org.apache.cassandra.io.sstable.SSTableIdentityIterator; import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.io.sstable.SSTableScanner; -import org.apache.cassandra.io.util.FileUtils; -import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.MergeIterator; -public class CompactionIterator extends ReducingIterator -implements Closeable, CompactionInfo.Holder +public class CompactionIterator extends AbstractIterator +implements CloseableIterator, CompactionInfo.Holder { private static Logger logger = LoggerFactory.getLogger(CompactionIterator.class); public static final int FILE_BUFFER_SIZE = 1024 * 1024; - protected final List rows = new ArrayList(); + private final MergeIterator source; protected final CompactionType type; protected final CompactionController controller; @@ -65,33 +67,26 @@ implements Closeable, CompactionInfo.Holder public CompactionIterator(CompactionType type, Iterable sstables, CompactionController controller) throws IOException { - this(type, getCollatingIterator(sstables), controller); + this(type, getScanners(sstables), controller); } - @SuppressWarnings("unchecked") - protected CompactionIterator(CompactionType type, Iterator iter, CompactionController controller) + protected CompactionIterator(CompactionType type, List scanners, CompactionController controller) { - super(iter); this.type = type; this.controller = controller; + this.source = MergeIterator.get(scanners, ICOMP, new Reducer()); row = 0; totalBytes = bytesRead = 0; - for (SSTableScanner scanner : getScanners()) - { + for (SSTableScanner scanner : scanners) totalBytes += scanner.getFileLength(); - } } - @SuppressWarnings("unchecked") - protected static CollatingIterator getCollatingIterator(Iterable sstables) throws IOException + protected static List getScanners(Iterable sstables) throws IOException { - // TODO CollatingIterator iter = FBUtilities.getCollatingIterator(); - CollatingIterator iter = FBUtilities.getCollatingIterator(); + ArrayList scanners = new ArrayList(); for (SSTableReader sstable : sstables) - { - iter.addIterator(sstable.getDirectScanner(FILE_BUFFER_SIZE)); - } - return iter; + scanners.add(sstable.getDirectScanner(FILE_BUFFER_SIZE)); + return scanners; } public CompactionInfo getCompactionInfo() @@ -103,50 +98,12 @@ implements Closeable, CompactionInfo.Holder totalBytes); } - @Override - protected boolean isEqual(SSTableIdentityIterator o1, SSTableIdentityIterator o2) + + public AbstractCompactedRow computeNext() { - return o1.getKey().equals(o2.getKey()); - } - - public void reduce(SSTableIdentityIterator current) - { - rows.add(current); - } - - protected AbstractCompactedRow getReduced() - { - assert rows.size() > 0; - - try - { - AbstractCompactedRow compactedRow = controller.getCompactedRow(rows); - if (compactedRow.isEmpty()) - { - controller.invalidateCachedRow(compactedRow.key); - return null; - } - - // If the raw is cached, we call removeDeleted on it to have/ coherent query returns. However it would look - // like some deleted columns lived longer than gc_grace + compaction. This can also free up big amount of - // memory on long running instances - controller.removeDeletedInCache(compactedRow.key); - - return compactedRow; - } - finally - { - rows.clear(); - if ((row++ % 1000) == 0) - { - bytesRead = 0; - for (SSTableScanner scanner : getScanners()) - { - bytesRead += scanner.getFilePointer(); - } - throttle(); - } - } + if (!source.hasNext()) + return endOfData(); + return source.next(); } private void throttle() @@ -187,16 +144,69 @@ implements Closeable, CompactionInfo.Holder public void close() throws IOException { - FileUtils.close(getScanners()); + source.close(); } protected Iterable getScanners() { - return ((CollatingIterator)source).getIterators(); + return (Iterable)(source.iterators()); } public String toString() { return this.getCompactionInfo().toString(); } + + protected class Reducer extends MergeIterator.Reducer + { + protected final List rows = new ArrayList(); + + public void reduce(IColumnIterator current) + { + rows.add((SSTableIdentityIterator)current); + } + + protected AbstractCompactedRow getReduced() + { + assert rows.size() > 0; + + try + { + AbstractCompactedRow compactedRow = controller.getCompactedRow(rows); + if (compactedRow.isEmpty()) + { + controller.invalidateCachedRow(compactedRow.key); + return null; + } + + // If the raw is cached, we call removeDeleted on it to have/ coherent query returns. However it would look + // like some deleted columns lived longer than gc_grace + compaction. This can also free up big amount of + // memory on long running instances + controller.removeDeletedInCache(compactedRow.key); + + return compactedRow; + } + finally + { + rows.clear(); + if ((row++ % 1000) == 0) + { + bytesRead = 0; + for (SSTableScanner scanner : getScanners()) + { + bytesRead += scanner.getFilePointer(); + } + throttle(); + } + } + } + } + + public final static Comparator ICOMP = new Comparator() + { + public int compare(IColumnIterator i1, IColumnIterator i2) + { + return i1.getKey().compareTo(i2.getKey()); + } + }; } diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 7019ca9576..6e6b3f127e 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -31,7 +31,6 @@ import javax.management.MBeanServer; import javax.management.ObjectName; import org.apache.commons.collections.PredicateUtils; -import org.apache.commons.collections.iterators.CollatingIterator; import org.apache.commons.collections.iterators.FilterIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -934,18 +933,16 @@ public class CompactionManager implements CompactionManagerMBean public ValidationCompactionIterator(ColumnFamilyStore cfs, Range range) throws IOException { super(CompactionType.VALIDATION, - getCollatingIterator(cfs.getSSTables(), range), + getScanners(cfs.getSSTables(), range), new CompactionController(cfs, cfs.getSSTables(), getDefaultGcBefore(cfs), true)); } - protected static CollatingIterator getCollatingIterator(Iterable sstables, Range range) throws IOException + protected static List getScanners(Iterable sstables, Range range) throws IOException { - CollatingIterator iter = FBUtilities.getCollatingIterator(); + ArrayList scanners = new ArrayList(); for (SSTableReader sstable : sstables) - { - iter.addIterator(sstable.getDirectScanner(FILE_BUFFER_SIZE, range)); - } - return iter; + scanners.add(sstable.getDirectScanner(FILE_BUFFER_SIZE, range)); + return scanners; } } diff --git a/src/java/org/apache/cassandra/db/compaction/LazilyCompactedRow.java b/src/java/org/apache/cassandra/db/compaction/LazilyCompactedRow.java index 557bd24c84..0f1a6a4c90 100644 --- a/src/java/org/apache/cassandra/db/compaction/LazilyCompactedRow.java +++ b/src/java/org/apache/cassandra/db/compaction/LazilyCompactedRow.java @@ -29,7 +29,6 @@ import java.util.*; import com.google.common.base.Predicates; import com.google.common.collect.Iterators; -import org.apache.commons.collections.iterators.CollatingIterator; import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.ColumnFamilyStore; @@ -40,7 +39,7 @@ import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.sstable.SSTableIdentityIterator; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.IIterableColumns; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.utils.MergeIterator; /** * LazilyCompactedRow only computes the row bloom filter and column index in memory @@ -60,7 +59,7 @@ public class LazilyCompactedRow extends AbstractCompactedRow implements IIterabl private final boolean shouldPurge; private final DataOutputBuffer headerBuffer; private ColumnFamily emptyColumnFamily; - private LazyColumnIterator iter; + private Reducer reducer; private int columnCount; private long columnSerializedSize; @@ -84,10 +83,10 @@ public class LazilyCompactedRow extends AbstractCompactedRow implements IIterabl // initialize row header so isEmpty can be called headerBuffer = new DataOutputBuffer(); ColumnIndexer.serialize(this, headerBuffer); - // reach into iterator used by ColumnIndexer to get column count and size - columnCount = iter.size; - columnSerializedSize = iter.serializedSize; - iter = null; + // reach into the reducer used during iteration to get column count and size + columnCount = reducer.size; + columnSerializedSize = reducer.serializedSize; + reducer = null; } public void write(DataOutput out) throws IOException @@ -156,10 +155,9 @@ public class LazilyCompactedRow extends AbstractCompactedRow implements IIterabl public Iterator iterator() { for (SSTableIdentityIterator row : rows) - { row.reset(); - } - iter = new LazyColumnIterator(new CollatingIterator(getComparator().columnComparator, rows)); + reducer = new Reducer(); + Iterator iter = MergeIterator.get(rows, getComparator().columnComparator, reducer); return Iterators.filter(iter, Predicates.notNull()); } @@ -168,23 +166,12 @@ public class LazilyCompactedRow extends AbstractCompactedRow implements IIterabl return columnCount; } - private class LazyColumnIterator extends ReducingIterator + private class Reducer extends MergeIterator.Reducer { ColumnFamily container = emptyColumnFamily.cloneMeShallow(); long serializedSize = 4; // int for column count int size = 0; - public LazyColumnIterator(Iterator source) - { - super(source); - } - - @Override - protected boolean isEqual(IColumn o1, IColumn o2) - { - return o1.name().equals(o2.name()); - } - public void reduce(IColumn current) { container.addColumn(current); diff --git a/src/java/org/apache/cassandra/db/filter/QueryFilter.java b/src/java/org/apache/cassandra/db/filter/QueryFilter.java index d4cccda4a5..27f824a51c 100644 --- a/src/java/org/apache/cassandra/db/filter/QueryFilter.java +++ b/src/java/org/apache/cassandra/db/filter/QueryFilter.java @@ -22,10 +22,7 @@ package org.apache.cassandra.db.filter; import java.nio.ByteBuffer; -import java.util.Comparator; -import java.util.Iterator; -import java.util.SortedSet; -import java.util.TreeSet; +import java.util.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -38,7 +35,8 @@ import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.io.util.FileDataInput; import org.apache.cassandra.thrift.SlicePredicate; import org.apache.cassandra.thrift.SliceRange; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.MergeIterator; public class QueryFilter { @@ -88,11 +86,14 @@ public class QueryFilter return superFilter.getSSTableColumnIterator(sstable, file, key); } - public void collectCollatedColumns(final ColumnFamily returnCF, Iterator collatedColumns, final int gcBefore) + // TODO move gcBefore into a field + public void collateColumns(final ColumnFamily returnCF, List> toCollate, AbstractType comparator, final int gcBefore) { + IFilter topLevelFilter = (superFilter == null ? filter : superFilter); + Comparator fcomp = topLevelFilter.getColumnComparator(comparator); // define a 'reduced' iterator that merges columns w/ the same name, which // greatly simplifies computing liveColumns in the presence of tombstones. - ReducingIterator reduced = new ReducingIterator(collatedColumns) + Iterator reduced = MergeIterator.get(toCollate, fcomp, new MergeIterator.Reducer() { ColumnFamily curCF = returnCF.cloneMeShallow(); @@ -137,9 +138,9 @@ public class QueryFilter return c; } - }; + }); - (superFilter == null ? filter : superFilter).collectReducedColumns(returnCF, reduced, gcBefore); + topLevelFilter.collectReducedColumns(returnCF, reduced, gcBefore); } public String getColumnFamilyName() diff --git a/src/java/org/apache/cassandra/io/sstable/KeyIterator.java b/src/java/org/apache/cassandra/io/sstable/KeyIterator.java index 4f3ca43ab0..b63a661d90 100644 --- a/src/java/org/apache/cassandra/io/sstable/KeyIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/KeyIterator.java @@ -21,7 +21,6 @@ package org.apache.cassandra.io.sstable; */ -import java.io.Closeable; import java.io.File; import java.io.IOError; import java.io.IOException; @@ -33,8 +32,9 @@ import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.CloseableIterator; -public class KeyIterator extends AbstractIterator implements Iterator, Closeable +public class KeyIterator extends AbstractIterator implements CloseableIterator { private final BufferedRandomAccessFile in; private final Descriptor desc; diff --git a/src/java/org/apache/cassandra/io/sstable/ReducingKeyIterator.java b/src/java/org/apache/cassandra/io/sstable/ReducingKeyIterator.java index 3f57bba8d7..42ed6cf7a6 100644 --- a/src/java/org/apache/cassandra/io/sstable/ReducingKeyIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/ReducingKeyIterator.java @@ -21,31 +21,26 @@ package org.apache.cassandra.io.sstable; */ -import java.io.Closeable; import java.io.IOException; +import java.util.ArrayList; import java.util.Collection; import java.util.Iterator; -import org.apache.commons.collections.iterators.CollatingIterator; - import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.utils.FBUtilities; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.MergeIterator; -public class ReducingKeyIterator implements Iterator, Closeable +public class ReducingKeyIterator implements CloseableIterator { - private final CollatingIterator ci; - private final ReducingIterator iter; + private final MergeIterator mi; public ReducingKeyIterator(Collection sstables) { - ci = FBUtilities.getCollatingIterator(); + ArrayList iters = new ArrayList(); for (SSTableReader sstable : sstables) - { - ci.addIterator(new KeyIterator(sstable.descriptor)); - } - - iter = new ReducingIterator(ci) + iters.add(new KeyIterator(sstable.descriptor)); + mi = MergeIterator.get(iters, DecoratedKey.comparator, new MergeIterator.Reducer() { DecoratedKey reduced = null; @@ -58,21 +53,21 @@ public class ReducingKeyIterator implements Iterator, Closeable { return reduced; } - }; + }); } public void close() throws IOException { - for (Object o : ci.getIterators()) + for (Object o : mi.iterators()) { - ((KeyIterator) o).close(); + ((CloseableIterator)o).close(); } } public long getTotalBytes() { long m = 0; - for (Object o : ci.getIterators()) + for (Object o : mi.iterators()) { m += ((KeyIterator) o).getTotalBytes(); } @@ -82,7 +77,7 @@ public class ReducingKeyIterator implements Iterator, Closeable public long getBytesRead() { long m = 0; - for (Object o : ci.getIterators()) + for (Object o : mi.iterators()) { m += ((KeyIterator) o).getBytesRead(); } @@ -96,12 +91,12 @@ public class ReducingKeyIterator implements Iterator, Closeable public boolean hasNext() { - return iter.hasNext(); + return mi.hasNext(); } public DecoratedKey next() { - return iter.next(); + return mi.next(); } public void remove() diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java index 93e07e26c7..7f768180d7 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java @@ -19,7 +19,6 @@ package org.apache.cassandra.io.sstable; -import java.io.Closeable; import java.io.File; import java.io.IOError; import java.io.IOException; @@ -34,9 +33,9 @@ import org.apache.cassandra.db.columniterator.IColumnIterator; import org.apache.cassandra.db.filter.QueryFilter; import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.CloseableIterator; - -public class SSTableScanner implements Iterator, Closeable +public class SSTableScanner implements CloseableIterator { private static Logger logger = LoggerFactory.getLogger(SSTableScanner.class); diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java index 7b6f6b9a0e..6db6d98518 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java +++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java @@ -24,7 +24,6 @@ import java.util.*; import java.util.concurrent.LinkedBlockingQueue; import com.google.common.collect.AbstractIterator; -import org.apache.commons.collections.iterators.CollatingIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -34,7 +33,8 @@ import org.apache.cassandra.db.RangeSliceReply; import org.apache.cassandra.db.Row; import org.apache.cassandra.net.Message; import org.apache.cassandra.utils.Pair; -import org.apache.cassandra.utils.ReducingIterator; +import org.apache.cassandra.utils.CloseableIterator; +import org.apache.cassandra.utils.MergeIterator; /** * Turns RangeSliceReply objects into row (string -> CF) maps, resolving @@ -64,35 +64,27 @@ public class RangeSliceResponseResolver implements IResponseResolver resolve() throws IOException { - CollatingIterator collator = new CollatingIterator(new Comparator>() - { - public int compare(Pair o1, Pair o2) - { - return o1.left.key.compareTo(o2.left.key); - } - }); - + ArrayList iters = new ArrayList(responses.size()); int n = 0; for (Message response : responses) { RangeSliceReply reply = RangeSliceReply.read(response.getMessageBody(), response.getVersion()); n = Math.max(n, reply.rows.size()); - collator.addIterator(new RowIterator(reply.rows.iterator(), response.getFrom())); + iters.add(new RowIterator(reply.rows.iterator(), response.getFrom())); } - // for each row, compute the combination of all different versions seen, and repair incomplete versions - return new ReducingIterator, Row>(collator) + MergeIterator, Row> iter = MergeIterator.get(iters, new Comparator>() + { + public int compare(Pair o1, Pair o2) + { + return o1.left.key.compareTo(o2.left.key); + } + }, new MergeIterator.Reducer, Row>() { List versions = new ArrayList(sources.size()); List versionSources = new ArrayList(sources.size()); DecoratedKey key; - @Override - protected boolean isEqual(Pair o1, Pair o2) - { - return o1.left.key.equals(o2.left.key); - } - public void reduce(Pair current) { key = current.left.key; @@ -122,7 +114,13 @@ public class RangeSliceResponseResolver implements IResponseResolver resolvedRows = new ArrayList(n); + while (iter.hasNext()) + resolvedRows.add(iter.next()); + + return resolvedRows; } public void preprocess(Message message) @@ -135,7 +133,7 @@ public class RangeSliceResponseResolver implements IResponseResolver> + private static class RowIterator extends AbstractIterator> implements CloseableIterator> { private final Iterator iter; private final InetAddress source; @@ -150,6 +148,8 @@ public class RangeSliceResponseResolver implements IResponseResolver(iter.next(), source) : endOfData(); } + + public void close() {} } public Iterable getMessages() diff --git a/src/java/org/apache/cassandra/service/RowRepairResolver.java b/src/java/org/apache/cassandra/service/RowRepairResolver.java index 727d44bffd..dc667c8882 100644 --- a/src/java/org/apache/cassandra/service/RowRepairResolver.java +++ b/src/java/org/apache/cassandra/service/RowRepairResolver.java @@ -26,8 +26,6 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; -import org.apache.commons.collections.iterators.CollatingIterator; - import org.apache.cassandra.db.*; import org.apache.cassandra.db.columniterator.IdentityQueryFilter; import org.apache.cassandra.db.filter.QueryFilter; @@ -35,6 +33,7 @@ import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.utils.*; public class RowRepairResolver extends AbstractRowResolver { @@ -142,14 +141,14 @@ public class RowRepairResolver extends AbstractRowResolver // this will handle removing columns and subcolumns that are supressed by a row or // supercolumn tombstone. QueryFilter filter = new QueryFilter(null, new QueryPath(resolved.metadata().cfName), new IdentityQueryFilter()); - CollatingIterator iter = new CollatingIterator(resolved.metadata().comparator.columnComparator); + List> iters = new ArrayList>(); for (ColumnFamily version : versions) { if (version == null) continue; - iter.addIterator(version.getColumnsMap().values().iterator()); + iters.add(FBUtilities.closeableIterator(version.getColumnsMap().values().iterator())); } - filter.collectCollatedColumns(resolved, iter, Integer.MIN_VALUE); + filter.collateColumns(resolved, iters, resolved.metadata().comparator, Integer.MIN_VALUE); return ColumnFamilyStore.removeDeleted(resolved, Integer.MIN_VALUE); } diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 6e1e619781..cd7a4c19e7 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -35,7 +35,7 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import com.google.common.base.Joiner; -import org.apache.commons.collections.iterators.CollatingIterator; +import com.google.common.collect.AbstractIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -331,32 +331,6 @@ public class FBUtilities } } - /* - TODO how to make this work w/ ReducingKeyIterator? - public static > CollatingIterator getCollatingIterator() - { - // CollatingIterator will happily NPE if you do not specify a comparator explicitly - return new CollatingIterator(new Comparator() - { - public int compare(T o1, T o2) - { - return o1.compareTo(o2); - } - }); - } - */ - public static CollatingIterator getCollatingIterator() - { - // CollatingIterator will happily NPE if you do not specify a comparator explicitly - return new CollatingIterator(new Comparator() - { - public int compare(Object o1, Object o2) - { - return ((Comparable) o1).compareTo(o2); - } - }); - } - public static void atomicSetMax(AtomicInteger atomic, int i) { while (true) @@ -614,4 +588,27 @@ public class FBUtilities return FBUtilities.construct(cache_provider, "row cache provider"); } + public static CloseableIterator closeableIterator(Iterator iterator) + { + return new WrappedCloseableIterator(iterator); + } + + private static final class WrappedCloseableIterator + extends AbstractIterator implements CloseableIterator + { + private final Iterator source; + public WrappedCloseableIterator(Iterator source) + { + this.source = source; + } + + protected T computeNext() + { + if (!source.hasNext()) + return endOfData(); + return source.next(); + } + + public void close() {} + } } diff --git a/src/java/org/apache/cassandra/utils/ReducingIterator.java b/src/java/org/apache/cassandra/utils/ReducingIterator.java deleted file mode 100644 index a96c2c9dcc..0000000000 --- a/src/java/org/apache/cassandra/utils/ReducingIterator.java +++ /dev/null @@ -1,87 +0,0 @@ -package org.apache.cassandra.utils; -/* - * - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - * - */ - - -import java.util.Iterator; - -import com.google.common.collect.AbstractIterator; - -/** - * reduces equal values from the source iterator to a single (optionally transformed) instance. - */ -public abstract class ReducingIterator extends AbstractIterator implements Iterator, Iterable -{ - protected Iterator source; - protected T1 last; - - public ReducingIterator(Iterator source) - { - this.source = source; - } - - /** combine this object with the previous ones. intermediate state is up to your implementation. */ - public abstract void reduce(T1 current); - - /** return the last object computed by reduce */ - protected abstract T2 getReduced(); - - /** override this if the keys you want to base the reduce on are not the same as the object itself (but can be generated from it) */ - protected boolean isEqual(T1 o1, T1 o2) - { - return o1.equals(o2); - } - - protected T2 computeNext() - { - if (last == null && !source.hasNext()) - return endOfData(); - - onKeyChange(); - boolean keyChanged = false; - while (!keyChanged) - { - if (last != null) - reduce(last); - if (!source.hasNext()) - { - last = null; - break; - } - T1 current = source.next(); - if (last != null && !isEqual(current, last)) - keyChanged = true; - last = current; - } - return getReduced(); - } - - /** - * Called at the begining of each new key, before any reduce is called. - * To be overriden by implementing classes. - */ - protected void onKeyChange() {} - - public Iterator iterator() - { - return this; - } -} diff --git a/test/unit/org/apache/cassandra/utils/MergeIteratorTest.java b/test/unit/org/apache/cassandra/utils/MergeIteratorTest.java new file mode 100644 index 0000000000..bd7535fb10 --- /dev/null +++ b/test/unit/org/apache/cassandra/utils/MergeIteratorTest.java @@ -0,0 +1,107 @@ +/* +* Licensed to the Apache Software Foundation (ASF) under one +* or more contributor license agreements. See the NOTICE file +* distributed with this work for additional information +* regarding copyright ownership. The ASF licenses this file +* to you under the Apache License, Version 2.0 (the +* "License"); you may not use this file except in compliance +* with the License. You may obtain a copy of the License at +* +* http://www.apache.org/licenses/LICENSE-2.0 +* +* Unless required by applicable law or agreed to in writing, +* software distributed under the License is distributed on an +* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +* KIND, either express or implied. See the License for the +* specific language governing permissions and limitations +* under the License. +*/ +package org.apache.cassandra.utils; + +import java.util.Arrays; +import java.util.Comparator; +import java.util.Iterator; +import java.util.List; + +import com.google.common.collect.AbstractIterator; +import com.google.common.collect.Iterators; +import com.google.common.collect.Ordering; + +import org.junit.Before; +import org.junit.Test; + +public class MergeIteratorTest +{ + CLI all = null, cat = null, a = null, b = null, c = null, d = null; + + @Before + public void clear() + { + all = new CLI("1", "2", "3", "3", "4", "5", "6", "7", "8", "8", "9"); + cat = new CLI("1", "2", "33", "4", "5", "6", "7", "88", "9"); + a = new CLI("1", "3", "5", "8"); + b = new CLI("2", "4", "6"); + c = new CLI("3", "7", "8", "9"); + d = new CLI(); + } + + @Test + public void testOneToOne() throws Exception + { + MergeIterator smi = MergeIterator.get(Arrays.asList(a, b, c, d), + Ordering.natural()); + assert Iterators.elementsEqual(all, smi); + smi.close(); + assert a.closed && b.closed && c.closed && d.closed; + } + + /** Test that duplicate values are concatted. */ + @Test + public void testManyToOne() throws Exception + { + MergeIterator.Reducer reducer = new MergeIterator.Reducer() + { + String concatted = ""; + public void reduce(String value) + { + concatted += value; + } + + public String getReduced() + { + String tmp = concatted; + concatted = ""; + return tmp; + } + }; + MergeIterator smi = MergeIterator.get(Arrays.asList(a, b, c, d), + Ordering.natural(), + reducer); + assert Iterators.elementsEqual(cat, smi); + smi.close(); + assert a.closed && b.closed && c.closed && d.closed; + } + + // closeable list iterator + public static class CLI extends AbstractIterator implements CloseableIterator + { + Iterator iter; + boolean closed = false; + public CLI(E... items) + { + this.iter = Arrays.asList(items).iterator(); + } + + protected E computeNext() + { + if (!iter.hasNext()) return endOfData(); + return iter.next(); + } + + public void close() + { + assert !this.closed; + this.closed = true; + } + } +}