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);
+        }
+    }
+}

Reply via email to