First draft of Collector utilities
Project: http://git-wip-us.apache.org/repos/asf/jena/repo Commit: http://git-wip-us.apache.org/repos/asf/jena/commit/e0c580c1 Tree: http://git-wip-us.apache.org/repos/asf/jena/tree/e0c580c1 Diff: http://git-wip-us.apache.org/repos/asf/jena/diff/e0c580c1 Branch: refs/heads/master Commit: e0c580c1241f874030f59ab28b3ba2a2b6f27775 Parents: 0ca10b7 Author: ajs6f <[email protected]> Authored: Sun Oct 1 10:02:28 2017 -0400 Committer: ajs6f <[email protected]> Committed: Fri Jan 5 09:26:07 2018 -0500 ---------------------------------------------------------------------- .../jena/query/util/DatasetCollector.java | 21 ++ .../query/util/DatasetIntoDatasetCollector.java | 42 ++++ .../org/apache/jena/query/util/DatasetLib.java | 21 ++ .../query/util/ModelIntoDatasetCollector.java | 53 ++++ .../jena/sparql/util/UnionDatasetGraph.java | 243 +++++++++++++++++++ .../jena/atlas/lib/IdentityFinishCollector.java | 13 + .../org/apache/jena/util/ModelCollector.java | 21 ++ .../jena/util/ModelIntoModelCollector.java | 37 +++ .../jena/util/StatementIntoModelCollector.java | 38 +++ 9 files changed, 489 insertions(+) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-arq/src/main/java/org/apache/jena/query/util/DatasetCollector.java ---------------------------------------------------------------------- diff --git a/jena-arq/src/main/java/org/apache/jena/query/util/DatasetCollector.java b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetCollector.java new file mode 100644 index 0000000..140c3fc --- /dev/null +++ b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetCollector.java @@ -0,0 +1,21 @@ +package org.apache.jena.query.util; + +import java.util.function.BinaryOperator; +import java.util.function.Supplier; + +import org.apache.jena.atlas.lib.IdentityFinishCollector; +import org.apache.jena.query.Dataset; +import org.apache.jena.query.DatasetFactory; + +public interface DatasetCollector<Input> extends IdentityFinishCollector<Input, Dataset> { + + @Override + default Supplier<Dataset> supplier() { + return DatasetFactory::createGeneral; + } + + @Override + default BinaryOperator<Dataset> combiner() { + return DatasetLib::union; + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-arq/src/main/java/org/apache/jena/query/util/DatasetIntoDatasetCollector.java ---------------------------------------------------------------------- diff --git a/jena-arq/src/main/java/org/apache/jena/query/util/DatasetIntoDatasetCollector.java b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetIntoDatasetCollector.java new file mode 100644 index 0000000..64af22d --- /dev/null +++ b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetIntoDatasetCollector.java @@ -0,0 +1,42 @@ +package org.apache.jena.query.util; + +import static java.util.stream.Collector.Characteristics.CONCURRENT; +import static java.util.stream.Collector.Characteristics.IDENTITY_FINISH; +import static java.util.stream.Collector.Characteristics.UNORDERED; +import static org.apache.jena.system.Txn.executeRead; +import static org.apache.jena.system.Txn.executeWrite; + +import java.util.Set; +import java.util.function.BiConsumer; + +import org.apache.jena.ext.com.google.common.collect.ImmutableSet; +import org.apache.jena.query.Dataset; + +public class DatasetIntoDatasetCollector implements DatasetCollector<Dataset> { + + @Override + public BiConsumer<Dataset, Dataset> accumulator() { + return (d1, d2) -> { + d1.getDefaultModel().add(d2.getDefaultModel()); + d2.listNames().forEachRemaining(name -> d1.getNamedModel(name).add(d2.getNamedModel(name))); + }; + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH); + } + + public static class ConcurrentStatementIntoModelCollector extends DatasetIntoDatasetCollector { + + @Override + public BiConsumer<Dataset, Dataset> accumulator() { + return (d1, d2) -> executeRead(d2, () -> executeWrite(d1, () -> super.accumulator().accept(d1, d2))); + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH, CONCURRENT); + } + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-arq/src/main/java/org/apache/jena/query/util/DatasetLib.java ---------------------------------------------------------------------- diff --git a/jena-arq/src/main/java/org/apache/jena/query/util/DatasetLib.java b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetLib.java new file mode 100644 index 0000000..8d1d789 --- /dev/null +++ b/jena-arq/src/main/java/org/apache/jena/query/util/DatasetLib.java @@ -0,0 +1,21 @@ +package org.apache.jena.query.util; + +import org.apache.jena.query.Dataset; + +public class DatasetLib { + + public static Dataset union(final Dataset d1, final Dataset d2) { + // TODO + throw new UnsupportedOperationException(); + } + + public static Dataset intersection(final Dataset d1, final Dataset d2) { + // TODO + throw new UnsupportedOperationException(); + } + + public static Dataset difference(final Dataset d1, final Dataset d2) { + // TODO + throw new UnsupportedOperationException(); + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-arq/src/main/java/org/apache/jena/query/util/ModelIntoDatasetCollector.java ---------------------------------------------------------------------- diff --git a/jena-arq/src/main/java/org/apache/jena/query/util/ModelIntoDatasetCollector.java b/jena-arq/src/main/java/org/apache/jena/query/util/ModelIntoDatasetCollector.java new file mode 100644 index 0000000..290c760 --- /dev/null +++ b/jena-arq/src/main/java/org/apache/jena/query/util/ModelIntoDatasetCollector.java @@ -0,0 +1,53 @@ +package org.apache.jena.query.util; + +import static java.util.stream.Collector.Characteristics.CONCURRENT; +import static java.util.stream.Collector.Characteristics.IDENTITY_FINISH; +import static java.util.stream.Collector.Characteristics.UNORDERED; +import static org.apache.jena.system.Txn.executeWrite; + +import java.util.Set; +import java.util.function.BiConsumer; + +import org.apache.jena.ext.com.google.common.collect.ImmutableSet; +import org.apache.jena.query.Dataset; +import org.apache.jena.rdf.model.Model; +import org.apache.jena.sparql.core.Quad; + +public class ModelIntoDatasetCollector implements DatasetCollector<Model> { + + private String graphName; + + public ModelIntoDatasetCollector(String graphName) { + this.graphName = graphName; + } + + /** + * Collects models into the default graph. + */ + public ModelIntoDatasetCollector() { + this(Quad.defaultGraphIRI.getURI()); + } + + @Override + public BiConsumer<Dataset, Model> accumulator() { + return (d, m) -> d.getNamedModel(graphName).add(m); + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH); + } + + public static class ConcurrentStatementIntoModelCollector extends ModelIntoDatasetCollector { + + @Override + public BiConsumer<Dataset, Model> accumulator() { + return (d, m) -> m.executeInTxn(() -> executeWrite(d, () -> super.accumulator().accept(d, m))); + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH, CONCURRENT); + } + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-arq/src/main/java/org/apache/jena/sparql/util/UnionDatasetGraph.java ---------------------------------------------------------------------- diff --git a/jena-arq/src/main/java/org/apache/jena/sparql/util/UnionDatasetGraph.java b/jena-arq/src/main/java/org/apache/jena/sparql/util/UnionDatasetGraph.java new file mode 100644 index 0000000..b857417 --- /dev/null +++ b/jena-arq/src/main/java/org/apache/jena/sparql/util/UnionDatasetGraph.java @@ -0,0 +1,243 @@ +package org.apache.jena.sparql.util; + +import static org.apache.jena.query.ReadWrite.WRITE; + +import java.util.Iterator; +import java.util.function.Function; + +import org.apache.jena.ext.com.google.common.collect.Iterators; +import org.apache.jena.graph.Graph; +import org.apache.jena.graph.Node; +import org.apache.jena.graph.compose.Union; +import org.apache.jena.query.ReadWrite; +import org.apache.jena.shared.Lock; +import org.apache.jena.sparql.core.DatasetGraph; +import org.apache.jena.sparql.core.Quad; + +public class UnionDatasetGraph implements DatasetGraph { + + private final DatasetGraph left, right; + + private final Lock lock; + + public UnionDatasetGraph(DatasetGraph left, DatasetGraph right) { + this.left = left; + this.right = right; + this.lock = new UnionLock(left.getLock(), right.getLock()); + } + + static Graph union(Graph left, Graph right) { + return new Union(left, right); + } + + Graph union(Function<DatasetGraph, Graph> op) { + return union(op.apply(left), op.apply(right)); + } + + boolean both(Function<DatasetGraph, Boolean> op) { + return op.apply(left) && op.apply(right); + } + + boolean either(Function<DatasetGraph, Boolean> op) { + return op.apply(left) || op.apply(right); + } + + <T> Iterator<T> fromEach(Function<DatasetGraph, Iterator<T>> op) { + return Iterators.concat(op.apply(left), op.apply(right)); + } + + @Override + public void begin(ReadWrite readWrite) { + if (readWrite.equals(WRITE)) throw new UnsupportedOperationException(); + left.begin(readWrite); + right.begin(readWrite); + } + + @Override + public void commit() { + throw new UnsupportedOperationException(); + } + + @Override + public void abort() { + throw new UnsupportedOperationException(); + } + + @Override + public void end() { + left.end(); + right.end(); + } + + @Override + public boolean isInTransaction() { + return either(DatasetGraph::isInTransaction); + } + + @Override + public Graph getDefaultGraph() { + return union(DatasetGraph::getDefaultGraph); + } + + @Override + public Graph getGraph(Node graphNode) { + return union(dsg -> dsg.getGraph(graphNode)); + } + + @Override + public Graph getUnionGraph() { + return union(DatasetGraph::getUnionGraph); + } + + @Override + public boolean containsGraph(Node graphNode) { + return either(dsg -> dsg.containsGraph(graphNode)); + } + + @Override + public void setDefaultGraph(Graph g) { + throw new UnsupportedOperationException(); + } + + @Override + public void addGraph(Node graphName, Graph graph) { + throw new UnsupportedOperationException(); + } + + @Override + public void removeGraph(Node graphName) { + throw new UnsupportedOperationException(); + } + + @Override + public Iterator<Node> listGraphNodes() { + return fromEach(DatasetGraph::listGraphNodes); + } + + @Override + public void add(Quad quad) { + throw new UnsupportedOperationException(); + } + + @Override + public void delete(Quad quad) { + throw new UnsupportedOperationException(); + } + + @Override + public void add(Node g, Node s, Node p, Node o) { + throw new UnsupportedOperationException(); + } + + @Override + public void delete(Node g, Node s, Node p, Node o) { + throw new UnsupportedOperationException(); + } + + @Override + public void deleteAny(Node g, Node s, Node p, Node o) { + throw new UnsupportedOperationException(); + } + + @Override + public Iterator<Quad> find() { + return fromEach(DatasetGraph::find); + } + + @Override + public Iterator<Quad> find(Quad quad) { + return fromEach(dsg -> dsg.find(quad)); + } + + @Override + public Iterator<Quad> find(Node g, Node s, Node p, Node o) { + return fromEach(dsg -> dsg.find(g, s, p, o)); + } + + @Override + public Iterator<Quad> findNG(Node g, Node s, Node p, Node o) { + return fromEach(dsg -> dsg.findNG(g, s, p, o)); + } + + @Override + public boolean contains(Node g, Node s, Node p, Node o) { + return either(dsg -> dsg.contains(g, s, p, o)); + } + + @Override + public boolean contains(Quad quad) { + return either(dsg -> dsg.contains(quad)); + } + + @Override + public void clear() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean isEmpty() { + return both(DatasetGraph::isEmpty); + } + + @Override + public Lock getLock() { + return lock; + } + + @Override + public Context getContext() { + // TODO Auto-generated method stub + return null; + } + + @Override + public long size() { + return left.size() + right.size(); + } + + @Override + public void close() { + throw new UnsupportedOperationException(); + } + + @Override + public boolean supportsTransactions() { + return both(DatasetGraph::supportsTransactions); + } + + @Override + public boolean supportsTransactionAbort() { + return both(DatasetGraph::supportsTransactionAbort); + } + + private static class UnionLock implements Lock { + + public UnionLock(Lock left, Lock right) { + this.left = left; + this.right = right; + } + + private final Lock left, right; + + @Override + public void enterCriticalSection(boolean readLockRequested) { + left.enterCriticalSection(readLockRequested); + right.enterCriticalSection(readLockRequested); + } + + @Override + public void leaveCriticalSection() { + left.leaveCriticalSection(); + right.leaveCriticalSection(); + } + } + + private static class UnionContext extends Context { + + UnionContext(Context left, Context right) { + this.context = left.context; + } + + + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-base/src/main/java/org/apache/jena/atlas/lib/IdentityFinishCollector.java ---------------------------------------------------------------------- diff --git a/jena-base/src/main/java/org/apache/jena/atlas/lib/IdentityFinishCollector.java b/jena-base/src/main/java/org/apache/jena/atlas/lib/IdentityFinishCollector.java new file mode 100644 index 0000000..472b73e --- /dev/null +++ b/jena-base/src/main/java/org/apache/jena/atlas/lib/IdentityFinishCollector.java @@ -0,0 +1,13 @@ +package org.apache.jena.atlas.lib; + +import java.util.function.Function; +import java.util.stream.Collector; + +public interface IdentityFinishCollector<T, A> extends Collector<T, A, A> { + + @Override + default Function<A, A> finisher() { + return Function.identity(); + } + +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-core/src/main/java/org/apache/jena/util/ModelCollector.java ---------------------------------------------------------------------- diff --git a/jena-core/src/main/java/org/apache/jena/util/ModelCollector.java b/jena-core/src/main/java/org/apache/jena/util/ModelCollector.java new file mode 100644 index 0000000..96a81e2 --- /dev/null +++ b/jena-core/src/main/java/org/apache/jena/util/ModelCollector.java @@ -0,0 +1,21 @@ +package org.apache.jena.util; + +import java.util.function.BinaryOperator; +import java.util.function.Supplier; + +import org.apache.jena.atlas.lib.IdentityFinishCollector; +import org.apache.jena.rdf.model.Model; +import org.apache.jena.rdf.model.ModelFactory; + +public interface ModelCollector<Input> extends IdentityFinishCollector<Input, Model> { + + @Override + default Supplier<Model> supplier() { + return ModelFactory::createDefaultModel; + } + + @Override + default BinaryOperator<Model> combiner() { + return ModelFactory::createUnion; + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-core/src/main/java/org/apache/jena/util/ModelIntoModelCollector.java ---------------------------------------------------------------------- diff --git a/jena-core/src/main/java/org/apache/jena/util/ModelIntoModelCollector.java b/jena-core/src/main/java/org/apache/jena/util/ModelIntoModelCollector.java new file mode 100644 index 0000000..413ae8f --- /dev/null +++ b/jena-core/src/main/java/org/apache/jena/util/ModelIntoModelCollector.java @@ -0,0 +1,37 @@ +package org.apache.jena.util; + +import static java.util.stream.Collector.Characteristics.CONCURRENT; +import static java.util.stream.Collector.Characteristics.IDENTITY_FINISH; +import static java.util.stream.Collector.Characteristics.UNORDERED; + +import java.util.Set; +import java.util.function.BiConsumer; + +import org.apache.jena.ext.com.google.common.collect.ImmutableSet; +import org.apache.jena.rdf.model.Model; + +public class ModelIntoModelCollector implements ModelCollector<Model> { + + @Override + public BiConsumer<Model, Model> accumulator() { + return Model::add; + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH); + } + + public static class ConcurrentModelIntoModelCollector extends ModelIntoModelCollector { + + @Override + public BiConsumer<Model, Model> accumulator() { + return (m1, m2) -> m1.executeInTxn(() -> m2.executeInTxn(() -> super.accumulator().accept(m1, m2))); + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH, CONCURRENT); + } + } +} http://git-wip-us.apache.org/repos/asf/jena/blob/e0c580c1/jena-core/src/main/java/org/apache/jena/util/StatementIntoModelCollector.java ---------------------------------------------------------------------- diff --git a/jena-core/src/main/java/org/apache/jena/util/StatementIntoModelCollector.java b/jena-core/src/main/java/org/apache/jena/util/StatementIntoModelCollector.java new file mode 100644 index 0000000..773bea5 --- /dev/null +++ b/jena-core/src/main/java/org/apache/jena/util/StatementIntoModelCollector.java @@ -0,0 +1,38 @@ +package org.apache.jena.util; + +import static java.util.stream.Collector.Characteristics.CONCURRENT; +import static java.util.stream.Collector.Characteristics.IDENTITY_FINISH; +import static java.util.stream.Collector.Characteristics.UNORDERED; + +import java.util.Set; +import java.util.function.BiConsumer; + +import org.apache.jena.ext.com.google.common.collect.ImmutableSet; +import org.apache.jena.rdf.model.Model; +import org.apache.jena.rdf.model.Statement; + +public class StatementIntoModelCollector implements ModelCollector<Statement> { + + @Override + public BiConsumer<Model, Statement> accumulator() { + return Model::add; + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH); + } + + public static class ConcurrentStatementIntoModelCollector extends StatementIntoModelCollector { + + @Override + public BiConsumer<Model, Statement> accumulator() { + return (m, s) -> m.executeInTxn(() -> m.add(s)); + } + + @Override + public Set<Characteristics> characteristics() { + return ImmutableSet.of(UNORDERED, IDENTITY_FINISH, CONCURRENT); + } + } +}
