Repository: incubator-blur Updated Branches: refs/heads/master 6b0b62d90 -> fc8bc600e
Streaming updates for large row is alomost complete, only thing left is to iterate through the row from the index and test that rows larger then can fit into memory can be updated. Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/fc8bc600 Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/fc8bc600 Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/fc8bc600 Branch: refs/heads/master Commit: fc8bc600ecc2759d80a0b5fc78791d56c0288ae9 Parents: 6b0b62d Author: Aaron McCurry <[email protected]> Authored: Tue Dec 9 17:33:05 2014 -0500 Committer: Aaron McCurry <[email protected]> Committed: Tue Dec 9 17:33:05 2014 -0500 ---------------------------------------------------------------------- .../blur/manager/writer/MutatableAction.java | 234 +++++++++++++------ .../writer/RecordToDocumentIterable.java | 1 - 2 files changed, 166 insertions(+), 69 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/fc8bc600/blur-core/src/main/java/org/apache/blur/manager/writer/MutatableAction.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/manager/writer/MutatableAction.java b/blur-core/src/main/java/org/apache/blur/manager/writer/MutatableAction.java index 5cd13fa..f933f5b 100644 --- a/blur-core/src/main/java/org/apache/blur/manager/writer/MutatableAction.java +++ b/blur-core/src/main/java/org/apache/blur/manager/writer/MutatableAction.java @@ -72,6 +72,64 @@ public class MutatableAction extends IndexAction { _writeRowMeter = Metrics.newMeter(metricName2, "Row/s", TimeUnit.SECONDS); } + static abstract class BaseRecordMutatorIterator implements Iterable<Record> { + + private final Record _record; + private final Iterable<Record> _iterable; + private final String _recordId; + + public BaseRecordMutatorIterator(Iterable<Record> iterable, Record record) { + _record = record; + _recordId = record.getRecordId(); + _iterable = iterable; + } + + protected abstract Record handleRecordMutate(Record existingRecord, Record newRecord); + + @Override + public Iterator<Record> iterator() { + final Iterator<Record> iterator = _iterable.iterator(); + return new Iterator<Record>() { + + private boolean _applied = false; + private boolean _append = false; + + @Override + public boolean hasNext() { + boolean hasNext = iterator.hasNext(); + if (hasNext) { + return true; + } + if (_applied) { + return false; // Already applied changes, finished. + } + _append = true; + return true; // Still need to add new record. + } + + @Override + public Record next() { + if (_append) { + _applied = true; + return _record; + } + Record record = iterator.next(); + if (record.getRecordId().equals(_recordId)) { + record = handleRecordMutate(record, _record); + _applied = true; + } + return record; + } + + @Override + public void remove() { + throw new RuntimeException("Not Supported."); + } + + }; + } + } + static class UpdateRow extends InternalAction { static abstract class UpdateRowAction { @@ -102,18 +160,54 @@ public class MutatableAction extends IndexAction { if (row == null) { return null; } else { - List<Record> records = new ArrayList<Record>(); - for (Record record : row) { - if (!record.getRecordId().equals(recordId)) { - records.add(record); - } - } - return new IterableRow(row.getRowId(), records); + return new IterableRow(row.getRowId(), new DeleteRecordIterator(row, recordId)); } } }); } + static class DeleteRecordIterator implements Iterable<Record> { + + private final String _recordIdToDelete; + private final Iterable<Record> _iterable; + + public DeleteRecordIterator(Iterable<Record> iterable, String recordIdToDelete) { + _recordIdToDelete = recordIdToDelete; + _iterable = iterable; + } + + @Override + public Iterator<Record> iterator() { + final GenericPeekableIterator<Record> iterator = GenericPeekableIterator.wrap(_iterable.iterator()); + return new Iterator<Record>() { + + @Override + public boolean hasNext() { + Record record = iterator.peek(); + if (record == null) { + return false; + } + if (record.getRecordId().equals(_recordIdToDelete)) { + iterator.next();// Eat the delete + return hasNext();// Move to the next record + } + return iterator.hasNext(); + } + + @Override + public Record next() { + return iterator.next(); + } + + @Override + public void remove() { + throw new RuntimeException("Not Supported."); + } + }; + } + + } + void appendColumns(final Record record) { _actions.add(new UpdateRowAction() { @Override @@ -121,28 +215,28 @@ public class MutatableAction extends IndexAction { if (row == null) { return new IterableRow(_rowId, Arrays.asList(record)); } else { - String recordId = record.getRecordId(); - List<Record> records = new ArrayList<Record>(); - boolean found = false; - for (Record r : row) { - if (!r.getRecordId().equals(recordId)) { - records.add(r); - } else { - found = true; - // Append columns - r.getColumns().addAll(record.getColumns()); - records.add(r); - } - } - if (!found) { - records.add(record); - } - return new IterableRow(_rowId, records); + return new IterableRow(row.getRowId(), new AppendColumnsIterator(row, record)); } } }); } + static class AppendColumnsIterator extends BaseRecordMutatorIterator { + + public AppendColumnsIterator(Iterable<Record> iterable, Record record) { + super(iterable, record); + } + + @Override + protected Record handleRecordMutate(Record existingRecord, Record newRecord) { + for (Column column : newRecord.getColumns()) { + existingRecord.addToColumns(column); + } + return existingRecord; + } + + } + void replaceColumns(final Record record) { _actions.add(new UpdateRowAction() { @Override @@ -150,28 +244,26 @@ public class MutatableAction extends IndexAction { if (row == null) { return new IterableRow(_rowId, Arrays.asList(record)); } else { - String recordId = record.getRecordId(); - boolean found = false; - List<Record> records = new ArrayList<Record>(); - for (Record r : row) { - if (!r.getRecordId().equals(recordId)) { - records.add(r); - } else { - found = true; - // Replace columns - records.add(replaceColumns(r, record)); - } - } - if (!found) { - records.add(record); - } - return new IterableRow(_rowId, records); + return new IterableRow(_rowId, new ReplaceColumnsIterator(row, record)); } } }); } - protected Record replaceColumns(Record existing, Record newRecord) { + static class ReplaceColumnsIterator extends BaseRecordMutatorIterator { + + public ReplaceColumnsIterator(Iterable<Record> iterable, Record record) { + super(iterable, record); + } + + @Override + protected Record handleRecordMutate(Record existingRecord, Record newRecord) { + return replaceColumns(existingRecord, newRecord); + } + + } + + protected static Record replaceColumns(Record existing, Record newRecord) { Map<String, List<Column>> existingColumns = getColumnMap(existing.getColumns()); Map<String, List<Column>> newColumns = getColumnMap(newRecord.getColumns()); existingColumns.putAll(newColumns); @@ -182,7 +274,7 @@ public class MutatableAction extends IndexAction { return record; } - private List<Column> toList(Collection<List<Column>> values) { + private static List<Column> toList(Collection<List<Column>> values) { ArrayList<Column> list = new ArrayList<Column>(); for (List<Column> v : values) { list.addAll(v); @@ -190,7 +282,7 @@ public class MutatableAction extends IndexAction { return list; } - private Map<String, List<Column>> getColumnMap(List<Column> columns) { + private static Map<String, List<Column>> getColumnMap(List<Column> columns) { Map<String, List<Column>> columnMap = new TreeMap<String, List<Column>>(); for (Column column : columns) { String name = column.getName(); @@ -209,23 +301,29 @@ public class MutatableAction extends IndexAction { @Override IterableRow performAction(IterableRow row) { if (row == null) { + // New Row return new IterableRow(_rowId, Arrays.asList(record)); } else { - List<Record> records = new ArrayList<Record>(); - String recordId = record.getRecordId(); - for (Record r : row) { - if (!r.getRecordId().equals(recordId)) { - records.add(r); - } - } - // Add replacement - records.add(record); - return new IterableRow(_rowId, records); + // Existing Row + return new IterableRow(_rowId, new ReplaceRecordIterator(row, record)); } } }); } + static class ReplaceRecordIterator extends BaseRecordMutatorIterator { + + public ReplaceRecordIterator(Iterable<Record> iterable, Record record) { + super(iterable, record); + } + + @Override + protected Record handleRecordMutate(Record existingRecord, Record newRecord) { + return newRecord; + } + + } + @Override void performAction(IndexSearcherClosable searcher, IndexWriter writer) throws IOException { Selector selector = new Selector(); @@ -246,22 +344,22 @@ public class MutatableAction extends IndexAction { iterableRow = action.performAction(iterableRow); } Term term = createRowId(_rowId); - // if (iterableRow != null) { - RecordToDocumentIterable docsToUpdate = new RecordToDocumentIterable(iterableRow, _fieldManager); - Iterator<Iterable<Field>> iterator = docsToUpdate.iterator(); - final GenericPeekableIterator<Iterable<Field>> gpi = GenericPeekableIterator.wrap(iterator); - if (gpi.peek() != null) { - writer.updateDocuments(term, wrapPrimeDoc(new Iterable<Iterable<Field>>() { - @Override - public Iterator<Iterable<Field>> iterator() { - return gpi; - } - })); - } else { - writer.deleteDocuments(term); + if (iterableRow != null) { + RecordToDocumentIterable docsToUpdate = new RecordToDocumentIterable(iterableRow, _fieldManager); + Iterator<Iterable<Field>> iterator = docsToUpdate.iterator(); + final GenericPeekableIterator<Iterable<Field>> gpi = GenericPeekableIterator.wrap(iterator); + if (gpi.peek() != null) { + writer.updateDocuments(term, wrapPrimeDoc(new Iterable<Iterable<Field>>() { + @Override + public Iterator<Iterable<Field>> iterator() { + return gpi; + } + })); + } else { + writer.deleteDocuments(term); + } + _writeRecordsMeter.mark(docsToUpdate.count()); } - _writeRecordsMeter.mark(docsToUpdate.count()); - // } _writeRowMeter.mark(); } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/fc8bc600/blur-core/src/main/java/org/apache/blur/manager/writer/RecordToDocumentIterable.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/manager/writer/RecordToDocumentIterable.java b/blur-core/src/main/java/org/apache/blur/manager/writer/RecordToDocumentIterable.java index d198938..a442d35 100644 --- a/blur-core/src/main/java/org/apache/blur/manager/writer/RecordToDocumentIterable.java +++ b/blur-core/src/main/java/org/apache/blur/manager/writer/RecordToDocumentIterable.java @@ -18,7 +18,6 @@ package org.apache.blur.manager.writer; import java.io.IOException; import java.util.Iterator; -import java.util.List; import org.apache.blur.analysis.FieldManager; import org.apache.blur.thrift.generated.Record;
