http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataContextTest.java ---------------------------------------------------------------------- diff --git a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataContextTest.java b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataContextTest.java new file mode 100644 index 0000000..1952640 --- /dev/null +++ b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataContextTest.java @@ -0,0 +1,613 @@ +/** + * 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.metamodel.mongodb.mongo3; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import com.mongodb.WriteConcern; + +import org.apache.metamodel.DataContext; +import org.apache.metamodel.UpdateCallback; +import org.apache.metamodel.UpdateScript; +import org.apache.metamodel.data.DataSet; +import org.apache.metamodel.data.InMemoryDataSet; +import org.apache.metamodel.query.FunctionType; +import org.apache.metamodel.query.SelectItem; +import org.apache.metamodel.schema.ColumnType; +import org.apache.metamodel.schema.Schema; +import org.apache.metamodel.schema.Table; +import org.apache.metamodel.util.SimpleTableDef; +import org.bson.Document; + +import com.mongodb.MongoClient; +import com.mongodb.MongoClientURI; +import com.mongodb.client.MongoCollection; +import com.mongodb.client.MongoDatabase; + +public class MongoDbDataContextTest extends MongoDbTestCase { + + private MongoDatabase mongoDb; + private MongoClient client; + + @Override + protected void setUp() throws Exception { + super.setUp(); + if (isConfigured()) { + client = new MongoClient(new MongoClientURI("mongodb://"+getHostname()+":27017")); + mongoDb = client.getDatabase(getDatabaseName()); + } + } + + @Override + protected void tearDown() throws Exception { + super.tearDown(); + if (isConfigured()) { + mongoDb.drop(); + client.close(); + } + } + + public void testNestedObjectFetching() throws Exception { + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + mongoDb.createCollection(getCollectionName()); + MongoCollection<Document> col = mongoDb.getCollection(getCollectionName()); + + final List<Document> list = new ArrayList<Document>(); + list.add(new Document().append("city", "Copenhagen").append("country", "Denmark")); + list.add(new Document().append("city", "Stockholm").append("country", "Sweden")); + + final Document dbRow = new Document(); + dbRow.append("name", new Document().append("first", "John").append("last", "Doe")); + dbRow.append("gender", "MALE"); + dbRow.append("addresses", list); + col.insertOne(dbRow); + + final MongoDbDataContext dc = new MongoDbDataContext(mongoDb, new SimpleTableDef(getCollectionName(), new String[] { + "name.first", "name.last", "gender", "addresses", "addresses[0].city", "addresses[0].country", + "addresses[5].foobar" })); + + final DataSet ds = dc.query().from(getCollectionName()).selectAll().execute(); + try { + assertTrue(ds.next()); + final Object addresses = ds.getRow().getValue(3); + assertEquals("Row[values=[John, Doe, MALE, " + addresses + ", Copenhagen, Denmark, null]]", ds.getRow() + .toString()); + assertTrue(addresses instanceof List); + assertFalse(ds.next()); + } finally { + ds.close(); + } + } + + public void testQueriesWithAutoGeneratedID() throws Exception { + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + + mongoDb.createCollection(getCollectionName()); + MongoCollection<Document> col = mongoDb.getCollection(getCollectionName()); + col.withWriteConcern(WriteConcern.ACKNOWLEDGED); + // create a couple of entries + + Document dbRow1 = new Document(); + dbRow1.put("name", "Mr. Black"); + dbRow1.put("category", "gen_id"); + dbRow1.put("age", 20); + col.insertOne(dbRow1); + final String autoGenID1 = dbRow1.get("_id").toString(); + + Document dbRow2 = new Document(); + dbRow2.put("name", "Mr. Pink"); + dbRow2.put("category", "gen_id"); + dbRow2.put("age", 40); + col.insertOne(dbRow2); + String autoGenID2 = dbRow2.get("_id").toString(); + + Document dbRow3 = new Document(); + dbRow3.put("_id", "123"); + dbRow3.put("name", "Mr. White"); + dbRow3.put("category", "gen_id"); + dbRow3.put("age", 30); + col.insertOne(dbRow3); + String fixedID3 = dbRow3.get("_id").toString(); + + final MongoDbDataContext dc = new MongoDbDataContext(mongoDb); + DataSet ds; + + // check all 3 entries inserted + ds = dc.query().from(getCollectionName()).selectAll() + .where("category").eq("gen_id").execute(); + assertEquals(3, ds.toRows().size()); + ds.close(); + + // select by autogenerated id + ds = dc.query().from(getCollectionName()).select("name").where("_id").eq(autoGenID1).execute(); + assertTrue(ds.next()); + assertEquals("Mr. Black", ds.getRow().getValue(0)); + ds.close(); + + // select by multiple autogenerated ids + ds = dc.query().from(getCollectionName()).select("name") + .where("_id").eq(autoGenID1) + .or("_id").eq(autoGenID2) + .execute(); + assertEquals(2, ds.toRows().size()); + ds.close(); + + // select by both autogenerated id and fixed id + ds = dc.query().from(getCollectionName()).select("name") + .where("_id").eq(autoGenID1) + .or("_id").eq(fixedID3) + .execute(); + assertEquals(2, ds.toRows().size()); + ds.close(); + + // delete by id + dc.executeUpdate(new UpdateScript() { + @Override + public void run(UpdateCallback callback) { + callback.deleteFrom(getCollectionName()) + .where("_id").eq(autoGenID1) + .execute(); + } + }); + + // select by autogenerated id which was deleted + ds = dc.query().from(getCollectionName()).select("name").where("_id").eq(autoGenID1).execute(); + assertEquals(0, ds.toRows().size()); + ds.close(); + + } + + public void testFirstRowAndMaxRows() throws Exception { + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + + mongoDb.createCollection(getCollectionName()); + MongoCollection<Document> col = mongoDb.getCollection(getCollectionName()); + + // create 3 records + for (int i = 0; i < 3; i++) { + Document dbRow = new Document(); + dbRow.put("id", i + 1); + col.insertOne(dbRow); + } + + final MongoDbDataContext dc = new MongoDbDataContext(mongoDb); + + DataSet ds; + + ds = dc.query().from(getCollectionName()).select("id").firstRow(2).execute(); + assertTrue(ds instanceof MongoDbDataSet); + assertTrue(ds.next()); + assertEquals("Row[values=[2]]", ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals("Row[values=[3]]", ds.getRow().toString()); + assertFalse(ds.next()); + ds.close(); + + ds = dc.query().from(getCollectionName()).select("id").maxRows(1).execute(); + assertTrue(ds instanceof MongoDbDataSet); + assertTrue(ds.next()); + assertEquals("Row[values=[1]]", ds.getRow().toString()); + assertFalse(ds.next()); + ds.close(); + + ds = dc.query().from(getCollectionName()).select("id").maxRows(1).firstRow(2).execute(); + assertTrue(ds instanceof MongoDbDataSet); + assertTrue(ds.next()); + assertEquals("Row[values=[2]]", ds.getRow().toString()); + assertFalse(ds.next()); + ds.close(); + } + + public void testRead() throws Exception { + // Adding a comment to commit something and invoke a build in Travis... + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + + mongoDb.createCollection(getCollectionName()); + MongoCollection<Document> col = mongoDb.getCollection(getCollectionName()); + + // create 1000 records + for (int i = 0; i < 1000; i++) { + Document dbRow = new Document(); + dbRow.put("id", i); + dbRow.put("name", "record no. " + i); + if (i % 5 == 0) { + dbRow.put("foo", "bar"); + } else { + dbRow.put("foo", "baz"); + } + Document nestedObj = new Document(); + nestedObj.put("count", i); + nestedObj.put("constant", "foobarbaz"); + dbRow.put("baz", nestedObj); + + dbRow.put("list", Arrays.<Object> asList("l1", "l2", "l3", i)); + + col.insertOne(dbRow); + } + + // Instantiate the actual data context + final DataContext dataContext = new MongoDbDataContext(mongoDb); + + assertTrue(Arrays.asList(dataContext.getDefaultSchema().getTableNames()).contains(getCollectionName())); + + Table table = dataContext.getDefaultSchema().getTableByName(getCollectionName()); + assertEquals("[_id, baz, foo, id, list, name]", Arrays.toString(table.getColumnNames())); + + assertEquals(ColumnType.MAP, table.getColumnByName("baz").getType()); + assertEquals(ColumnType.VARCHAR, table.getColumnByName("foo").getType()); + assertEquals(ColumnType.LIST, table.getColumnByName("list").getType()); + assertEquals(ColumnType.INTEGER, table.getColumnByName("id").getType()); + assertEquals(ColumnType.ROWID, table.getColumnByName("_id").getType()); + + DataSet ds = dataContext.query().from(getCollectionName()).select("name").and("foo").and("baz").and("list") + .where("id").greaterThan(800).or("foo").isEquals("bar").execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + try { + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 0, bar, Document{{count=0, constant=foobarbaz}}, [l1, l2, l3, 0]]]", + ds.getRow().toString()); + + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 5, bar, Document{{count=5, constant=foobarbaz}}, [l1, l2, l3, 5]]]", + ds.getRow().toString()); + + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 10, bar, Document{{count=10, constant=foobarbaz}}, [l1, l2, l3, 10]]]", + ds.getRow().toString()); + + for (int j = 15; j < 801; j++) { + if (j % 5 == 0) { + assertTrue(ds.next()); + assertEquals("Row[values=[record no. " + j + ", bar, Document{{count=" + j + + ", constant=foobarbaz}}, [l1, l2, l3, " + j + "]]]", ds.getRow() + .toString()); + } + } + + assertTrue(ds.next()); + assertTrue(ds.getRow().getValue(2) instanceof Map); + assertEquals(Document.class, ds.getRow().getValue(2).getClass()); + + assertTrue("unexpected type: " + ds.getRow().getValue(3).getClass(), + ds.getRow().getValue(3) instanceof List); + assertEquals(ArrayList.class, ds.getRow().getValue(3).getClass()); + + assertEquals( + "Row[values=[record no. 801, baz, Document{{count=801, constant=foobarbaz}}, [l1, l2, l3, 801]]]", + ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 802, baz, Document{{count=802, constant=foobarbaz}}, [l1, l2, l3, 802]]]", + ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 803, baz, Document{{count=803, constant=foobarbaz}}, [l1, l2, l3, 803]]]", + ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 804, baz, Document{{count=804, constant=foobarbaz}}, [l1, l2, l3, 804]]]", + ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals( + "Row[values=[record no. 805, bar, Document{{count=805, constant=foobarbaz}}, [l1, l2, l3, 805]]]", + ds.getRow().toString()); + + for (int i = 0; i < 194; i++) { + assertTrue(ds.next()); + } + assertEquals( + "Row[values=[record no. 999, baz, Document{{count=999, constant=foobarbaz}}, [l1, l2, l3, 999]]]", + ds.getRow().toString()); + assertFalse(ds.next()); + } finally { + ds.close(); + } + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("id").in(2, 6, 8, 9) + .execute(); + assertTrue(ds.next()); + assertEquals("Row[values=[2, record no. 2]]", ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals("Row[values=[6, record no. 6]]", ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals("Row[values=[8, record no. 8]]", ds.getRow().toString()); + assertTrue(ds.next()); + assertEquals("Row[values=[9, record no. 9]]", ds.getRow().toString()); + assertFalse(ds.next()); + ds.close(); + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("foo").isEquals("bar") + .execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + + try { + List<Object[]> objectArrays = ds.toObjectArrays(); + assertEquals(200, objectArrays.size()); + assertEquals("[0, record no. 0]", Arrays.toString(objectArrays.get(0))); + } finally { + ds.close(); + } + + // test GREATER_THAN_OR_EQUAL + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("id") + .greaterThanOrEquals(500).and("foo").isEquals("bar").execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + + try { + List<Object[]> objectArrays = ds.toObjectArrays(); + assertEquals(100, objectArrays.size()); + assertEquals("[500, record no. 500]", Arrays.toString(objectArrays.get(0))); + } finally { + ds.close(); + } + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("id") + .greaterThanOrEquals(501).and("foo").isEquals("bar").execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + + try { + List<Object[]> objectArrays = ds.toObjectArrays(); + assertEquals(99, objectArrays.size()); + assertEquals("[505, record no. 505]", Arrays.toString(objectArrays.get(0))); + } finally { + ds.close(); + } + + // test LESS_THAN_OR_EQUAL + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("id").lessThanOrEquals(500) + .and("foo").isEquals("bar").execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + + try { + List<Object[]> objectArrays = ds.toObjectArrays(); + assertEquals(101, objectArrays.size()); + assertEquals("[500, record no. 500]", Arrays.toString(objectArrays.get(100))); + } finally { + ds.close(); + } + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("id").lessThanOrEquals(499) + .and("foo").isEquals("bar").execute(); + assertEquals(MongoDbDataSet.class, ds.getClass()); + assertFalse(((MongoDbDataSet) ds).isQueryPostProcessed()); + + try { + List<Object[]> objectArrays = ds.toObjectArrays(); + assertEquals(100, objectArrays.size()); + assertEquals("[495, record no. 495]", Arrays.toString(objectArrays.get(99))); + } finally { + ds.close(); + } + + // test a primary key lookup query + Document dbRow = new Document(); + dbRow.put("_id", 123456); + dbRow.put("id", 123456); + dbRow.put("name", "record no. " + 123456); + dbRow.put("foo", "bar123456"); + Document nestedObj = new Document(); + nestedObj.put("count", 123456); + nestedObj.put("constant", "foobarbaz"); + dbRow.put("baz", nestedObj); + + dbRow.put("list", Arrays.<Object> asList("l1", "l2", "l3", 123456)); + + col.insertOne(dbRow); + + ds = dataContext.query().from(getCollectionName()).select("id").and("name").where("_id").eq(123456).execute(); + assertTrue(ds.next()); + assertEquals("Row[values=[123456, record no. 123456]]", ds.getRow().toString()); + assertFalse(ds.next()); + + // do a query that we cannot push to mongo + // Replace column index 0 by 1 + ds = dataContext.query().from(getCollectionName()) + .select(FunctionType.SUM, dataContext.getDefaultSchema().getTables()[0].getColumnByName("id")) + .where("foo").isEquals("bar").execute(); + assertEquals(InMemoryDataSet.class, ds.getClass()); + + ds.close(); + } + + public void testCreateAndWriteData() throws Exception { + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + final MongoDbDataContext dc = new MongoDbDataContext(mongoDb); + final Schema defaultSchema = dc.getDefaultSchema(); + + dc.executeUpdate(new UpdateScript() { + @Override + public void run(UpdateCallback callback) { + Table[] tables = defaultSchema.getTables(); + for (Table table : tables) { + callback.deleteFrom(table).execute(); + } + } + }); + + assertEquals(0, defaultSchema.getTableCount()); + + dc.executeUpdate(new UpdateScript() { + + @Override + public void run(UpdateCallback callback) { + Table table = callback.createTable(defaultSchema, "some_entries").withColumn("foo").withColumn("bar") + .withColumn("baz").withColumn("list").execute(); + + callback.insertInto(table).value("foo", 1).value("bar", "hello").execute(); + callback.insertInto(table).value("foo", 2).value("bar", "world").execute(); + callback.insertInto(table).value("foo", 3).value("bar", "hi").execute(); + + Map<String, Object> nestedObj = new HashMap<String, Object>(); + nestedObj.put("foo", "bar"); + nestedObj.put("123", 456); + + callback.insertInto(table).value("foo", 4).value("bar", "there").value("baz", nestedObj) + .value("list", Arrays.asList(1, 2, 3)).execute(); + } + }); + + DataSet dataSet; + assertEquals(1, defaultSchema.getTableCount()); + + // "Pure" SELECT COUNT(*) query + dataSet = dc.query().from("some_entries").selectCount().execute(); + dataSet.close(); + assertTrue(dataSet.next()); + assertEquals(1, dataSet.getSelectItems().length); + assertEquals(SelectItem.getCountAllItem(), dataSet.getSelectItems()[0]); + assertEquals(4l, dataSet.getRow().getValue(SelectItem.getCountAllItem())); + assertFalse(dataSet.next()); + assertEquals(InMemoryDataSet.class, dataSet.getClass()); + + // A conditional SELECT COUNT(*) query + dataSet = dc.query().from("some_entries").selectCount().where("foo").greaterThan(2).execute(); + dataSet.close(); + assertTrue(dataSet.next()); + assertEquals(1, dataSet.getSelectItems().length); + assertEquals(SelectItem.getCountAllItem(), dataSet.getSelectItems()[0]); + assertEquals(2l, dataSet.getRow().getValue(SelectItem.getCountAllItem())); + assertFalse(dataSet.next()); + assertEquals(InMemoryDataSet.class, dataSet.getClass()); + + // Select columns + dataSet = dc.query().from("some_entries").select("foo").and("bar").and("baz").and("list").execute(); + assertTrue(dataSet.next()); + assertEquals("Row[values=[1, hello, null, null]]", dataSet.getRow().toString()); + assertTrue(dataSet.next()); + assertEquals("Row[values=[2, world, null, null]]", dataSet.getRow().toString()); + assertTrue(dataSet.next()); + assertEquals("Row[values=[3, hi, null, null]]", dataSet.getRow().toString()); + assertTrue(dataSet.next()); + assertEquals("Row[values=[4, there, Document{{123=456, foo=bar}}, [1, 2, 3]]]", dataSet.getRow().toString()); + assertFalse(dataSet.next()); + dataSet.close(); + assertEquals(MongoDbDataSet.class, dataSet.getClass()); + + // delete some records + dc.executeUpdate(new UpdateScript() { + @Override + public void run(UpdateCallback callback) { + callback.deleteFrom("some_entries").where("foo").greaterThan(2).where("baz").isNotNull().execute(); + } + }); + + dataSet = dc.query().from("some_entries").select("foo").execute(); + assertTrue(dataSet.next()); + assertEquals("Row[values=[1]]", dataSet.getRow().toString()); + assertTrue(dataSet.next()); + assertEquals("Row[values=[2]]", dataSet.getRow().toString()); + assertTrue(dataSet.next()); + assertEquals("Row[values=[3]]", dataSet.getRow().toString()); + assertFalse(dataSet.next()); + dataSet.close(); + assertEquals(MongoDbDataSet.class, dataSet.getClass()); + + // drop the collection + dc.executeUpdate(new UpdateScript() { + @Override + public void run(UpdateCallback callback) { + callback.dropTable("some_entries").execute(); + } + }); + + assertNull(dc.getTableByQualifiedLabel("some_entries")); + + dc.refreshSchemas(); + assertEquals(0, defaultSchema.getTableCount()); + } + + public void testSelectWithLikeOperator() throws Exception { + if (!isConfigured()) { + System.err.println(getInvalidConfigurationMessage()); + return; + } + + mongoDb.createCollection(getCollectionName()); + MongoCollection<Document> col = mongoDb.getCollection(getCollectionName()); + + final Document dbRow = new Document(); + dbRow.append("name", new Document().append("first", "John").append("last", "Doe")); + dbRow.append("gender", "MALE"); + col.insertOne(dbRow); + + final Document dbRow2 = new Document(); + dbRow2.append("name", new Document().append("first", "Mary").append("last", "Johnson")); + dbRow2.append("gender", "FEMALE"); + col.insertOne(dbRow2); + + final Document dbRow3 = new Document(); + dbRow3.append("name", new Document().append("first", "X").append("last", "Unknown")); + dbRow3.append("gender", "UNKNOWN"); + col.insertOne(dbRow3); + + final MongoDbDataContext dc = new MongoDbDataContext(mongoDb, new SimpleTableDef(getCollectionName(), new String[] { + "name.first", "name.last", "gender", "addresses", "addresses[0].city", "addresses[0].country", + "addresses[5].foobar" })); + + final DataSet ds1 = dc.executeQuery("select * from my_collection where gender LIKE '%MALE%'"); + final DataSet ds2 = dc.executeQuery("select * from my_collection where gender LIKE 'MALE%'"); + final DataSet ds3 = dc.executeQuery("select * from my_collection where gender LIKE '%NK%OW%'"); + final DataSet ds4 = dc.executeQuery("select * from my_collection where gender LIKE '%MALE'"); + try { + assertTrue(ds1.next()); + assertTrue(ds1.next()); + assertFalse(ds1.next()); + assertTrue(ds2.next()); + assertFalse(ds2.next()); + assertTrue(ds3.next()); + assertFalse(ds3.next()); + assertTrue(ds4.next()); + assertTrue(ds4.next()); + assertFalse(ds4.next()); + } finally { + ds1.close(); + ds2.close(); + ds3.close(); + ds4.close(); + } + } +}
http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataCopyer.java ---------------------------------------------------------------------- diff --git a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataCopyer.java b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataCopyer.java new file mode 100644 index 0000000..42ef536 --- /dev/null +++ b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbDataCopyer.java @@ -0,0 +1,127 @@ +/** + * 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.metamodel.mongodb.mongo3; + +import java.io.File; +import java.sql.Connection; +import java.sql.DriverManager; + +import org.apache.metamodel.DataContext; +import org.apache.metamodel.UpdateCallback; +import org.apache.metamodel.UpdateScript; +import org.apache.metamodel.data.DataSet; +import org.apache.metamodel.data.Row; +import org.apache.metamodel.insert.RowInsertionBuilder; +import org.apache.metamodel.jdbc.JdbcDataContext; +import org.apache.metamodel.schema.Column; +import org.apache.metamodel.schema.Schema; +import org.apache.metamodel.schema.Table; +import org.apache.metamodel.util.FileHelper; + +import com.mongodb.MongoClient; +import com.mongodb.client.MongoDatabase; + + +/** + * Simple example program that can copy data to a MongoDB collection + */ +public class MongoDbDataCopyer { + + private final DataContext _sourceDataContext; + private final MongoDatabase _mongoDb; + private final String _collectionName; + private final String _sourceSchemaName; + private final String _sourceTableName; + + // example copy job that will populate the mongodb with Derby data + public static void main(String[] args) throws Exception { + System.setProperty("derby.storage.tempDirector", FileHelper.getTempDir().getAbsolutePath()); + System.setProperty("derby.stream.error.file", File.createTempFile("metamodel-derby", ".log").getAbsolutePath()); + + File dbFile = new File("../jdbc/src/test/resources/derby_testdb.jar"); + dbFile = dbFile.getCanonicalFile(); + if (!dbFile.exists()) { + throw new IllegalStateException("File does not exist: " + dbFile); + } + + Class.forName("org.apache.derby.jdbc.EmbeddedDriver"); + Connection connection = DriverManager.getConnection("jdbc:derby:jar:(" + dbFile.getAbsolutePath() + + ")derby_testdb;territory=en"); + connection.setReadOnly(true); + + MongoClient client = new MongoClient(); + MongoDatabase mongoDb = client.getDatabase("orderdb_copy"); + + DataContext sourceDataContext = new JdbcDataContext(connection); + + new MongoDbDataCopyer(mongoDb, "orders", sourceDataContext, "APP", "orders").copy(); + new MongoDbDataCopyer(mongoDb, "offices", sourceDataContext, "APP", "offices").copy(); + new MongoDbDataCopyer(mongoDb, "payments", sourceDataContext, "APP", "payments").copy(); + new MongoDbDataCopyer(mongoDb, "orderfact", sourceDataContext, "APP", "orderfact").copy(); + new MongoDbDataCopyer(mongoDb, "products", sourceDataContext, "APP", "products").copy(); + + connection.close(); + client.close(); + } + + public MongoDbDataCopyer(MongoDatabase mongoDb, String collectionName, DataContext sourceDataContext, String sourceSchemaName, + String sourceTableName) { + _mongoDb = mongoDb; + _collectionName = collectionName; + _sourceDataContext = sourceDataContext; + _sourceSchemaName = sourceSchemaName; + _sourceTableName = sourceTableName; + } + + public void copy() { + final MongoDbDataContext targetDataContext = new MongoDbDataContext(_mongoDb); + targetDataContext.executeUpdate(new UpdateScript() { + + @Override + public void run(UpdateCallback callback) { + final Table sourceTable = getSourceTable(); + final Table targetTable = callback.createTable(targetDataContext.getDefaultSchema(), _collectionName) + .like(sourceTable).execute(); + final Column[] sourceColumns = sourceTable.getColumns(); + final DataSet dataSet = _sourceDataContext.query().from(sourceTable).select(sourceColumns).execute(); + while (dataSet.next()) { + final Row row = dataSet.getRow(); + + RowInsertionBuilder insertBuilder = callback.insertInto(targetTable); + for (Column column : sourceColumns) { + insertBuilder = insertBuilder.value(column.getName(), row.getValue(column)); + } + insertBuilder.execute(); + } + dataSet.close(); + } + }); + } + + private Table getSourceTable() { + final Schema schema; + if (_sourceSchemaName != null) { + schema = _sourceDataContext.getSchemaByName(_sourceSchemaName); + } else { + schema = _sourceDataContext.getDefaultSchema(); + } + + return schema.getTableByName(_sourceTableName); + } +} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbTestCase.java ---------------------------------------------------------------------- diff --git a/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbTestCase.java b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbTestCase.java new file mode 100644 index 0000000..edc98c9 --- /dev/null +++ b/mongodb/mongo3/src/test/java/org/apache/metamodel/mongodb/mongo3/MongoDbTestCase.java @@ -0,0 +1,111 @@ +/** + * 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.metamodel.mongodb.mongo3; + +import java.io.File; +import java.io.FileNotFoundException; +import java.io.FileReader; +import java.io.IOException; +import java.util.Properties; + +import junit.framework.TestCase; + +public abstract class MongoDbTestCase extends TestCase { + + private static final String DEFAULT_TEST_COLLECTION_NAME = "my_collection"; + private static final String DEFAULT_TEST_DATABASE_NAME = "metamodel_test"; + + private String _hostname; + private String _collectionName; + private boolean _configured; + + private String _databaseName; + + @Override + protected void setUp() throws Exception { + super.setUp(); + + File file = new File(getPropertyFilePath()); + if (file.exists()) { + loadPropertyFile(file); + } else { + // Continuous integration case + if (System.getenv("CONTINUOUS_INTEGRATION") != null) { + File travisFile = new File("../travis-metamodel-integrationtest-configuration.properties"); + if (travisFile.exists()) { + loadPropertyFile(travisFile); + } else { + _configured = false; + } + } else { + _configured = false; + } + } + } + + private void loadPropertyFile(File file) throws FileNotFoundException, IOException { + Properties properties = new Properties(); + properties.load(new FileReader(file)); + _hostname = properties.getProperty("mongodb.hostname"); + + _databaseName = properties.getProperty("mongodb.databaseName"); + if (_databaseName == null || _databaseName.isEmpty()) { + _databaseName = DEFAULT_TEST_DATABASE_NAME; + } + + _collectionName = properties.getProperty("mongodb.collectionName"); + if (_collectionName == null || _collectionName.isEmpty()) { + _collectionName = DEFAULT_TEST_COLLECTION_NAME; + } + + _configured = (_hostname != null && !_hostname.isEmpty()); + + if (_configured) { + System.out.println("Loaded MongoDB configuration. Hostname=" + _hostname + ", Collection=" + + _collectionName); + } + + } + + private String getPropertyFilePath() { + String userHome = System.getProperty("user.home"); + return userHome + "/metamodel-integrationtest-configuration.properties"; + } + + protected String getInvalidConfigurationMessage() { + return "!!! WARN !!! MongoDB module ignored\r\n" + "Please configure mongodb connection locally (" + + getPropertyFilePath() + "), to run integration tests"; + } + + public String getDatabaseName() { + return _databaseName; + } + + public String getCollectionName() { + return _collectionName; + } + + public String getHostname() { + return _hostname; + } + + public boolean isConfigured() { + return _configured; + } +} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/pom.xml ---------------------------------------------------------------------- diff --git a/mongodb/pom.xml b/mongodb/pom.xml index 61a3784..c0a9448 100644 --- a/mongodb/pom.xml +++ b/mongodb/pom.xml @@ -25,41 +25,12 @@ under the License. </parent> <modelVersion>4.0.0</modelVersion> <artifactId>MetaModel-mongodb</artifactId> + <packaging>pom</packaging> <name>MetaModel module for MongoDB databases</name> - <dependencies> - <dependency> - <groupId>org.apache.metamodel</groupId> - <artifactId>MetaModel-core</artifactId> - <version>${project.version}</version> - </dependency> - <dependency> - <groupId>org.mongodb</groupId> - <artifactId>mongo-java-driver</artifactId> - <version>3.1.0</version> - </dependency> - - <!-- Test dependencies --> - <dependency> - <groupId>org.apache.metamodel</groupId> - <artifactId>MetaModel-jdbc</artifactId> - <version>${project.version}</version> - <scope>test</scope> - </dependency> - <dependency> - <groupId>org.apache.derby</groupId> - <artifactId>derby</artifactId> - <version>10.8.1.2</version> - <scope>test</scope> - </dependency> - <dependency> - <groupId>org.slf4j</groupId> - <artifactId>slf4j-nop</artifactId> - <scope>test</scope> - </dependency> - <dependency> - <groupId>junit</groupId> - <artifactId>junit</artifactId> - <scope>test</scope> - </dependency> - </dependencies> + + <modules> + <module>common</module> + <module>mongo2</module> + <module>mongo3</module> + </modules> </project> http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/DefaultWriteConcernAdvisor.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/DefaultWriteConcernAdvisor.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/DefaultWriteConcernAdvisor.java deleted file mode 100644 index 6ea1903..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/DefaultWriteConcernAdvisor.java +++ /dev/null @@ -1,32 +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.metamodel.mongodb; - -import com.mongodb.WriteConcern; - -/** - * Default {@link WriteConcernAdvisor}. Always returns - * {@link WriteConcern#NORMAL}. - */ -public class DefaultWriteConcernAdvisor extends SimpleWriteConcernAdvisor { - - public DefaultWriteConcernAdvisor() { - super(WriteConcern.NORMAL); - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDBUtils.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDBUtils.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDBUtils.java deleted file mode 100644 index cb69a04..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDBUtils.java +++ /dev/null @@ -1,77 +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.metamodel.mongodb; - -import java.util.List; -import java.util.Map; - -import org.apache.metamodel.data.DataSetHeader; -import org.apache.metamodel.data.DefaultRow; -import org.apache.metamodel.data.Row; -import org.apache.metamodel.query.SelectItem; -import org.apache.metamodel.schema.Column; -import org.apache.metamodel.util.CollectionUtils; - -import com.mongodb.DBObject; - -/** - * A utility class for MongoDB module. - */ -public class MongoDBUtils { - - /** - * Converts a MongoDB data object {@link DBObject} into MetaModel - * {@link Row}. - * - * @param dbObject - * a MongoDB object storing data. - * @param header - * a header describing the columns of the data stored. - * @return the MetaModel {@link Row} result object. - */ - public static Row toRow(DBObject dbObject, DataSetHeader header) { - if (dbObject == null) { - return null; - } - - final Map<?,?> map = dbObject.toMap(); - - final int size = header.size(); - final Object[] values = new Object[size]; - for (int i = 0; i < values.length; i++) { - final SelectItem selectItem = header.getSelectItem(i); - final String key = selectItem.getColumn().getName(); - final Object value = CollectionUtils.find(map, key); - values[i] = toValue(selectItem.getColumn(), value); - } - return new DefaultRow(header, values); - } - - private static Object toValue(Column column, Object value) { - if (value instanceof List) { - return value; - } - if (value instanceof DBObject) { - DBObject basicDBObject = (DBObject) value; - return basicDBObject.toMap(); - } - return value; - } - -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataContext.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataContext.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataContext.java deleted file mode 100644 index 26db6d4..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataContext.java +++ /dev/null @@ -1,526 +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.metamodel.mongodb; - -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Map.Entry; -import java.util.Set; -import java.util.SortedMap; -import java.util.TreeMap; -import java.util.regex.Pattern; - -import org.apache.metamodel.DataContext; -import org.apache.metamodel.MetaModelException; -import org.apache.metamodel.QueryPostprocessDataContext; -import org.apache.metamodel.UpdateScript; -import org.apache.metamodel.UpdateableDataContext; -import org.apache.metamodel.data.DataSet; -import org.apache.metamodel.data.DataSetHeader; -import org.apache.metamodel.data.InMemoryDataSet; -import org.apache.metamodel.data.Row; -import org.apache.metamodel.data.SimpleDataSetHeader; -import org.apache.metamodel.query.FilterItem; -import org.apache.metamodel.query.FromItem; -import org.apache.metamodel.query.OperatorType; -import org.apache.metamodel.query.Query; -import org.apache.metamodel.query.SelectItem; -import org.apache.metamodel.schema.Column; -import org.apache.metamodel.schema.ColumnType; -import org.apache.metamodel.schema.ColumnTypeImpl; -import org.apache.metamodel.schema.MutableColumn; -import org.apache.metamodel.schema.MutableSchema; -import org.apache.metamodel.schema.MutableTable; -import org.apache.metamodel.schema.Schema; -import org.apache.metamodel.schema.Table; -import org.apache.metamodel.util.SimpleTableDef; -import org.bson.types.ObjectId; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.mongodb.BasicDBList; -import com.mongodb.BasicDBObject; -import com.mongodb.DB; -import com.mongodb.DBCollection; -import com.mongodb.DBCursor; -import com.mongodb.DBObject; -import com.mongodb.WriteConcern; - -/** - * DataContext implementation for MongoDB. - * - * Since MongoDB has no schema, a virtual schema will be used in this - * DataContext. This implementation supports either automatic discovery of a - * schema or manual specification of a schema, through the - * {@link MongoDbTableDef} class. - */ -public class MongoDbDataContext extends QueryPostprocessDataContext implements UpdateableDataContext { - - private static final Logger logger = LoggerFactory.getLogger(MongoDbDataSet.class); - - private final DB _mongoDb; - private final SimpleTableDef[] _tableDefs; - private WriteConcernAdvisor _writeConcernAdvisor; - private Schema _schema; - - /** - * Constructor available for backwards compatibility - * - * @deprecated use {@link #MongoDbDataContext(DB, SimpleTableDef...)} - * instead - */ - @Deprecated - public MongoDbDataContext(DB mongoDb, MongoDbTableDef... tableDefs) { - this(mongoDb, (SimpleTableDef[]) tableDefs); - } - - /** - * Constructs a {@link MongoDbDataContext}. This constructor accepts a - * custom array of {@link MongoDbTableDef}s which allows the user to define - * his own view on the collections in the database. - * - * @param mongoDb - * the mongo db connection - * @param tableDefs - * an array of {@link MongoDbTableDef}s, which define the table - * and column model of the mongo db collections. (consider using - * {@link #detectSchema(DB)} or {@link #detectTable(DB, String)} - * ). - */ - public MongoDbDataContext(DB mongoDb, SimpleTableDef... tableDefs) { - _mongoDb = mongoDb; - _tableDefs = tableDefs; - _schema = null; - } - - /** - * Constructs a {@link MongoDbDataContext} and automatically detects the - * schema structure/view on all collections (see {@link #detectSchema(DB)}). - * - * @param mongoDb - * the mongo db connection - */ - public MongoDbDataContext(DB mongoDb) { - this(mongoDb, detectSchema(mongoDb)); - } - - /** - * Performs an analysis of the available collections in a Mongo {@link DB} - * instance and tries to detect the table's structure based on the first - * 1000 documents in each collection. - * - * @param db - * the mongo db to inspect - * @return a mutable schema instance, useful for further fine tuning by the - * user. - * @see #detectTable(DB, String) - */ - public static SimpleTableDef[] detectSchema(DB db) { - Set<String> collectionNames = db.getCollectionNames(); - SimpleTableDef[] result = new SimpleTableDef[collectionNames.size()]; - int i = 0; - for (String collectionName : collectionNames) { - SimpleTableDef table = detectTable(db, collectionName); - result[i] = table; - i++; - } - return result; - } - - /** - * Performs an analysis of an available collection in a Mongo {@link DB} - * instance and tries to detect the table structure based on the first 1000 - * documents in the collection. - * - * @param db - * the mongo DB - * @param collectionName - * the name of the collection - * @return a table definition for mongo db. - */ - public static SimpleTableDef detectTable(DB db, String collectionName) { - final DBCollection collection = db.getCollection(collectionName); - final DBCursor cursor = collection.find().limit(1000); - - final SortedMap<String, Set<Class<?>>> columnsAndTypes = new TreeMap<String, Set<Class<?>>>(); - while (cursor.hasNext()) { - DBObject object = cursor.next(); - Set<String> keysInObject = object.keySet(); - for (String key : keysInObject) { - Set<Class<?>> types = columnsAndTypes.get(key); - if (types == null) { - types = new HashSet<Class<?>>(); - columnsAndTypes.put(key, types); - } - Object value = object.get(key); - if (value != null) { - types.add(value.getClass()); - } - } - } - cursor.close(); - - final String[] columnNames = new String[columnsAndTypes.size()]; - final ColumnType[] columnTypes = new ColumnType[columnsAndTypes.size()]; - - int i = 0; - for (Entry<String, Set<Class<?>>> columnAndTypes : columnsAndTypes.entrySet()) { - final String columnName = columnAndTypes.getKey(); - final Set<Class<?>> columnTypeSet = columnAndTypes.getValue(); - final Class<?> columnType; - if (columnTypeSet.size() == 1) { - columnType = columnTypeSet.iterator().next(); - } else { - columnType = Object.class; - } - columnNames[i] = columnName; - if (columnType == ObjectId.class) { - columnTypes[i] = ColumnType.ROWID; - } else { - columnTypes[i] = ColumnTypeImpl.convertColumnType(columnType); - } - i++; - } - - return new SimpleTableDef(collectionName, columnNames, columnTypes); - } - - @Override - protected Schema getMainSchema() throws MetaModelException { - if (_schema == null) { - MutableSchema schema = new MutableSchema(getMainSchemaName()); - for (SimpleTableDef tableDef : _tableDefs) { - - MutableTable table = tableDef.toTable().setSchema(schema); - Column[] rowIdColumns = table.getColumnsOfType(ColumnType.ROWID); - for (Column column : rowIdColumns) { - if (column instanceof MutableColumn) { - ((MutableColumn) column).setPrimaryKey(true); - } - } - - schema.addTable(table); - } - - _schema = schema; - } - return _schema; - } - - @Override - protected String getMainSchemaName() throws MetaModelException { - return _mongoDb.getName(); - } - - @Override - protected Number executeCountQuery(Table table, List<FilterItem> whereItems, boolean functionApproximationAllowed) { - final DBCollection collection = _mongoDb.getCollection(table.getName()); - - final DBObject query = createMongoDbQuery(table, whereItems); - - logger.info("Executing MongoDB 'count' query: {}", query); - final long count = collection.count(query); - - return count; - } - - @Override - protected Row executePrimaryKeyLookupQuery(Table table, List<SelectItem> selectItems, Column primaryKeyColumn, - Object keyValue) { - final DBCollection collection = _mongoDb.getCollection(table.getName()); - - List<FilterItem> whereItems = new ArrayList<FilterItem>(); - SelectItem selectItem = new SelectItem(primaryKeyColumn); - FilterItem primaryKeyWhereItem = new FilterItem(selectItem, OperatorType.EQUALS_TO, keyValue); - whereItems.add(primaryKeyWhereItem); - final DBObject query = createMongoDbQuery(table, whereItems); - final DBObject resultDBObject = collection.findOne(query); - - DataSetHeader header = new SimpleDataSetHeader(selectItems); - - Row row = MongoDBUtils.toRow(resultDBObject, header); - - return row; - } - - @Override - public DataSet executeQuery(Query query) { - // Check for queries containing only simple selects and where clauses, - // or if it is a COUNT(*) query. - - // if from clause only contains a main schema table - List<FromItem> fromItems = query.getFromClause().getItems(); - if (fromItems.size() == 1 && fromItems.get(0).getTable() != null - && fromItems.get(0).getTable().getSchema() == _schema) { - final Table table = fromItems.get(0).getTable(); - - // if GROUP BY, HAVING and ORDER BY clauses are not specified - if (query.getGroupByClause().isEmpty() && query.getHavingClause().isEmpty() - && query.getOrderByClause().isEmpty()) { - - final List<FilterItem> whereItems = query.getWhereClause().getItems(); - - // if all of the select items are "pure" column selection - boolean allSelectItemsAreColumns = true; - List<SelectItem> selectItems = query.getSelectClause().getItems(); - - // if it is a - // "SELECT [columns] FROM [table] WHERE [conditions]" - // query. - for (SelectItem selectItem : selectItems) { - if (selectItem.getFunction() != null || selectItem.getColumn() == null) { - allSelectItemsAreColumns = false; - break; - } - } - - if (allSelectItemsAreColumns) { - logger.debug("Query can be expressed in full MongoDB, no post processing needed."); - - // prepare for a non-post-processed query - Column[] columns = new Column[selectItems.size()]; - for (int i = 0; i < columns.length; i++) { - columns[i] = selectItems.get(i).getColumn(); - } - - // checking if the query is a primary key lookup query - if (whereItems.size() == 1) { - final FilterItem whereItem = whereItems.get(0); - final SelectItem selectItem = whereItem.getSelectItem(); - if (!whereItem.isCompoundFilter() && selectItem != null && selectItem.getColumn() != null) { - final Column column = selectItem.getColumn(); - if (column.isPrimaryKey() && OperatorType.EQUALS_TO.equals(whereItem.getOperator())) { - logger.debug("Query is a primary key lookup query. Trying executePrimaryKeyLookupQuery(...)"); - final Object operand = whereItem.getOperand(); - final Row row = executePrimaryKeyLookupQuery(table, selectItems, column, operand); - if (row == null) { - logger.debug("DataContext did not return any primary key lookup query results. Proceeding " - + "with manual lookup."); - } else { - final DataSetHeader header = new SimpleDataSetHeader(selectItems); - return new InMemoryDataSet(header, row); - } - } - } - } - - int firstRow = (query.getFirstRow() == null ? 1 : query.getFirstRow()); - int maxRows = (query.getMaxRows() == null ? -1 : query.getMaxRows()); - - final DataSet dataSet = materializeMainSchemaTableInternal(table, columns, whereItems, firstRow, - maxRows, false); - return dataSet; - } - } - } - - logger.debug("Query will be simplified for MongoDB and post processed."); - return super.executeQuery(query); - } - - private DataSet materializeMainSchemaTableInternal(Table table, Column[] columns, List<FilterItem> whereItems, - int firstRow, int maxRows, boolean queryPostProcessed) { - final DBCollection collection = _mongoDb.getCollection(table.getName()); - - final DBObject query = createMongoDbQuery(table, whereItems); - - logger.info("Executing MongoDB 'find' query: {}", query); - DBCursor cursor = collection.find(query); - - if (maxRows > 0) { - cursor = cursor.limit(maxRows); - } - if (firstRow > 1) { - final int skip = firstRow - 1; - cursor = cursor.skip(skip); - } - - return new MongoDbDataSet(cursor, columns, queryPostProcessed); - } - - protected BasicDBObject createMongoDbQuery(Table table, List<FilterItem> whereItems) { - assert _schema == table.getSchema(); - - final BasicDBObject query = new BasicDBObject(); - if (whereItems != null && !whereItems.isEmpty()) { - for (FilterItem item : whereItems) { - convertToCursorObject(query, item); - } - } - - return query; - } - - private void convertToCursorObject(BasicDBObject query, FilterItem item) { - if (item.isCompoundFilter()) { - - BasicDBList orList = new BasicDBList(); - - final FilterItem[] childItems = item.getChildItems(); - for (FilterItem childItem : childItems) { - BasicDBObject childObject = new BasicDBObject(); - convertToCursorObject(childObject, childItem); - orList.add(childObject); - } - - query.put("$or", orList); - - } else { - - final Column column = item.getSelectItem().getColumn(); - final String columnName = column.getName(); - final String operatorName = getOperatorName(item); - - Object operand = item.getOperand(); - if (ObjectId.isValid(String.valueOf(operand))) { - operand = new ObjectId(String.valueOf(operand)); - } - - final BasicDBObject existingFilterObject = (BasicDBObject) query.get(columnName); - if (existingFilterObject == null) { - if (operatorName == null) { - if (OperatorType.LIKE.equals(item.getOperator())) { - query.put(columnName, turnOperandIntoRegExp(operand)); - } else { - query.put(columnName, operand); - } - } else { - query.put(columnName, new BasicDBObject(operatorName, operand)); - } - } else { - if (operatorName == null) { - throw new IllegalStateException("Cannot retrieve records for a column with two EQUALS_TO operators"); - } else { - existingFilterObject.append(operatorName, operand); - } - } - } - } - - private String getOperatorName(FilterItem item) { - final OperatorType operator = item.getOperator(); - - if (OperatorType.EQUALS_TO.equals(operator)) { - return null; - } - if (OperatorType.LIKE.equals(operator)) { - return null; - } - if (OperatorType.LESS_THAN.equals(operator)) { - return "$lt"; - } - if (OperatorType.LESS_THAN_OR_EQUAL.equals(operator)) { - return "$lte"; - } - if (OperatorType.GREATER_THAN.equals(operator)) { - return "$gt"; - } - if (OperatorType.GREATER_THAN_OR_EQUAL.equals(operator)) { - return "$gte"; - } - if (OperatorType.DIFFERENT_FROM.equals(operator)) { - return "$ne"; - } - if (OperatorType.IN.equals(operator)) { - return "$in"; - } - - throw new IllegalStateException("Unsupported operator type: " + operator); - } - - private Pattern turnOperandIntoRegExp(Object operand) { - StringBuilder operandAsRegExp = new StringBuilder(replaceWildCardLikeChars(operand.toString())); - operandAsRegExp.insert(0, "^").append("$"); - return Pattern.compile(operandAsRegExp.toString(), Pattern.CASE_INSENSITIVE); - } - - private String replaceWildCardLikeChars(String operand) { - return operand.replaceAll("%", ".*"); - } - - @Override - protected DataSet materializeMainSchemaTable(Table table, Column[] columns, int maxRows) { - return materializeMainSchemaTableInternal(table, columns, null, 1, maxRows, true); - } - - @Override - protected DataSet materializeMainSchemaTable(Table table, Column[] columns, int firstRow, int maxRows) { - return materializeMainSchemaTableInternal(table, columns, null, firstRow, maxRows, true); - } - - /** - * Executes an update with a specific {@link WriteConcernAdvisor}. - */ - public void executeUpdate(UpdateScript update, WriteConcernAdvisor writeConcernAdvisor) { - MongoDbUpdateCallback callback = new MongoDbUpdateCallback(this, writeConcernAdvisor); - try { - update.run(callback); - } finally { - callback.close(); - } - } - - /** - * Executes an update with a specific {@link WriteConcern}. - */ - public void executeUpdate(UpdateScript update, WriteConcern writeConcern) { - executeUpdate(update, new SimpleWriteConcernAdvisor(writeConcern)); - } - - @Override - public void executeUpdate(UpdateScript update) { - executeUpdate(update, getWriteConcernAdvisor()); - } - - /** - * Gets the {@link WriteConcernAdvisor} to use on - * {@link #executeUpdate(UpdateScript)} calls. - */ - public WriteConcernAdvisor getWriteConcernAdvisor() { - if (_writeConcernAdvisor == null) { - return new DefaultWriteConcernAdvisor(); - } - return _writeConcernAdvisor; - } - - /** - * Sets a global {@link WriteConcern} advisor to use on - * {@link #executeUpdate(UpdateScript)}. - */ - public void setWriteConcernAdvisor(WriteConcernAdvisor writeConcernAdvisor) { - _writeConcernAdvisor = writeConcernAdvisor; - } - - /** - * Gets the {@link DB} instance that this {@link DataContext} is backed by. - */ - public DB getMongoDb() { - return _mongoDb; - } - - protected void addTable(MutableTable table) { - if (_schema instanceof MutableSchema) { - MutableSchema mutableSchema = (MutableSchema) _schema; - mutableSchema.addTable(table); - } else { - throw new UnsupportedOperationException("Schema is not mutable"); - } - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataSet.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataSet.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataSet.java deleted file mode 100644 index ee522b8..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDataSet.java +++ /dev/null @@ -1,83 +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.metamodel.mongodb; - -import org.apache.metamodel.data.AbstractDataSet; -import org.apache.metamodel.data.Row; -import org.apache.metamodel.schema.Column; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.mongodb.DBCursor; -import com.mongodb.DBObject; - -final class MongoDbDataSet extends AbstractDataSet { - - private static final Logger logger = LoggerFactory.getLogger(MongoDbDataSet.class); - - private final DBCursor _cursor; - private final boolean _queryPostProcessed; - - private boolean _closed; - private volatile DBObject _dbObject; - - public MongoDbDataSet(DBCursor cursor, Column[] columns, boolean queryPostProcessed) { - super(columns); - _cursor = cursor; - _queryPostProcessed = queryPostProcessed; - _closed = false; - } - - public boolean isQueryPostProcessed() { - return _queryPostProcessed; - } - - @Override - public void close() { - super.close(); - _cursor.close(); - _closed = true; - } - - @Override - protected void finalize() throws Throwable { - super.finalize(); - if (!_closed) { - logger.warn("finalize() invoked, but DataSet is not closed. Invoking close() on {}", this); - close(); - } - } - - @Override - public boolean next() { - if (_cursor.hasNext()) { - _dbObject = _cursor.next(); - return true; - } else { - _dbObject = null; - return false; - } - } - - @Override - public Row getRow() { - return MongoDBUtils.toRow(_dbObject, getHeader()); - } - -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDeleteBuilder.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDeleteBuilder.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDeleteBuilder.java deleted file mode 100644 index f1b8591..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDeleteBuilder.java +++ /dev/null @@ -1,56 +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.metamodel.mongodb; - -import org.apache.metamodel.MetaModelException; -import org.apache.metamodel.delete.AbstractRowDeletionBuilder; -import org.apache.metamodel.schema.Table; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.mongodb.BasicDBObject; -import com.mongodb.DBCollection; -import com.mongodb.WriteConcern; -import com.mongodb.WriteResult; - -final class MongoDbDeleteBuilder extends AbstractRowDeletionBuilder { - - private static final Logger logger = LoggerFactory.getLogger(MongoDbDeleteBuilder.class); - - private final MongoDbUpdateCallback _updateCallback; - - public MongoDbDeleteBuilder(MongoDbUpdateCallback updateCallback, Table table) { - super(table); - _updateCallback = updateCallback; - } - - @Override - public void execute() throws MetaModelException { - final DBCollection collection = _updateCallback.getCollection(getTable().getName()); - - final MongoDbDataContext dataContext = _updateCallback.getDataContext(); - final BasicDBObject query = dataContext.createMongoDbQuery(getTable(), getWhereItems()); - - WriteConcern writeConcern = _updateCallback.getWriteConcernAdvisor().adviceDeleteQuery(collection, query); - - final WriteResult writeResult = collection.remove(query, writeConcern); - logger.info("Remove returned result: {}", writeResult); - } - -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDropTableBuilder.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDropTableBuilder.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDropTableBuilder.java deleted file mode 100644 index 4164346..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbDropTableBuilder.java +++ /dev/null @@ -1,43 +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.metamodel.mongodb; - -import org.apache.metamodel.MetaModelException; -import org.apache.metamodel.drop.AbstractTableDropBuilder; -import org.apache.metamodel.schema.MutableSchema; -import org.apache.metamodel.schema.Table; - -final class MongoDbDropTableBuilder extends AbstractTableDropBuilder { - - private final MongoDbUpdateCallback _updateCallback; - - public MongoDbDropTableBuilder(MongoDbUpdateCallback updateCallback, Table table) { - super(table); - _updateCallback = updateCallback; - } - - @Override - public void execute() throws MetaModelException { - Table table = getTable(); - _updateCallback.removeCollection(table.getName()); - MutableSchema schema = (MutableSchema) table.getSchema(); - schema.removeTable(table); - } - -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbInsertionBuilder.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbInsertionBuilder.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbInsertionBuilder.java deleted file mode 100644 index 412612c..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbInsertionBuilder.java +++ /dev/null @@ -1,64 +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.metamodel.mongodb; - -import org.apache.metamodel.MetaModelException; -import org.apache.metamodel.insert.AbstractRowInsertionBuilder; -import org.apache.metamodel.insert.RowInsertionBuilder; -import org.apache.metamodel.schema.Column; -import org.apache.metamodel.schema.Table; -import org.slf4j.LoggerFactory; -import org.slf4j.Logger; - -import com.mongodb.BasicDBObject; -import com.mongodb.DBCollection; -import com.mongodb.WriteConcern; -import com.mongodb.WriteResult; - -final class MongoDbInsertionBuilder extends AbstractRowInsertionBuilder<MongoDbUpdateCallback> implements RowInsertionBuilder { - - private static final Logger logger = LoggerFactory.getLogger(MongoDbInsertionBuilder.class); - - public MongoDbInsertionBuilder(MongoDbUpdateCallback updateCallback, Table table) { - super(updateCallback, table); - } - - @Override - public void execute() throws MetaModelException { - final Column[] columns = getColumns(); - final Object[] values = getValues(); - - final BasicDBObject doc = new BasicDBObject(); - - for (int i = 0; i < values.length; i++) { - Object value = values[i]; - if (value != null) { - doc.put(columns[i].getName(), value); - } - } - - final MongoDbUpdateCallback updateCallback = getUpdateCallback(); - final DBCollection collection = updateCallback.getCollection(getTable().getName()); - - final WriteConcern writeConcern = updateCallback.getWriteConcernAdvisor().adviceInsert(collection, doc); - - final WriteResult writeResult = collection.insert(doc, writeConcern); - logger.info("Insert returned result: {}", writeResult); - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableCreationBuilder.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableCreationBuilder.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableCreationBuilder.java deleted file mode 100644 index a57590d..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableCreationBuilder.java +++ /dev/null @@ -1,57 +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.metamodel.mongodb; - -import org.apache.metamodel.MetaModelException; -import org.apache.metamodel.create.AbstractTableCreationBuilder; -import org.apache.metamodel.create.TableCreationBuilder; -import org.apache.metamodel.schema.ColumnType; -import org.apache.metamodel.schema.ImmutableColumn; -import org.apache.metamodel.schema.MutableTable; -import org.apache.metamodel.schema.Schema; -import org.apache.metamodel.schema.Table; - -final class MongoDbTableCreationBuilder extends - AbstractTableCreationBuilder<MongoDbUpdateCallback> implements - TableCreationBuilder { - - public MongoDbTableCreationBuilder(MongoDbUpdateCallback updateCallback, - Schema schema, String name) { - super(updateCallback, schema, name); - } - - @Override - public Table execute() throws MetaModelException { - final MongoDbDataContext dataContext = getUpdateCallback() - .getDataContext(); - final Schema schema = dataContext.getDefaultSchema(); - final MutableTable table = getTable(); - if (table.getColumnByName("_id") == null) { - // all mongo db collections have an _id field as the first field. - ImmutableColumn idColumn = new ImmutableColumn("_id", - ColumnType.ROWID, table, table.getColumnCount(), null, - null, null, null, true, null, true); - table.addColumn(idColumn); - } - table.setSchema(schema); - getUpdateCallback().createCollection(table.getName()); - dataContext.addTable(table); - return table; - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableDef.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableDef.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableDef.java deleted file mode 100644 index 0a40fa1..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbTableDef.java +++ /dev/null @@ -1,46 +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.metamodel.mongodb; - -import java.io.Serializable; - -import org.apache.metamodel.schema.ColumnType; -import org.apache.metamodel.util.SimpleTableDef; - -/** - * Defines a table layout for {@link MongoDbDataContext} tables. This class can - * be used as an instruction set for the {@link MongoDbDataContext} to specify - * which collections, which columns (and their types) should be included in the - * schema structure of a Mongo DB database. - * - * @deprecated use {@link SimpleTableDef} instead. - */ -@Deprecated -public final class MongoDbTableDef extends SimpleTableDef implements Serializable { - - private static final long serialVersionUID = 1L; - - public MongoDbTableDef(String name, String[] columnNames, ColumnType[] columnTypes) { - super(name, columnNames, columnTypes); - } - - public MongoDbTableDef(String name, String[] columnNames) { - super(name, columnNames); - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbUpdateCallback.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbUpdateCallback.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbUpdateCallback.java deleted file mode 100644 index 0558898..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/MongoDbUpdateCallback.java +++ /dev/null @@ -1,115 +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.metamodel.mongodb; - -import java.io.Closeable; -import java.util.HashMap; -import java.util.Map; - -import org.apache.metamodel.AbstractUpdateCallback; -import org.apache.metamodel.UpdateCallback; -import org.apache.metamodel.create.TableCreationBuilder; -import org.apache.metamodel.delete.RowDeletionBuilder; -import org.apache.metamodel.drop.TableDropBuilder; -import org.apache.metamodel.insert.RowInsertionBuilder; -import org.apache.metamodel.schema.Schema; -import org.apache.metamodel.schema.Table; - -import com.mongodb.BasicDBObject; -import com.mongodb.DBCollection; - -final class MongoDbUpdateCallback extends AbstractUpdateCallback implements UpdateCallback, Closeable { - - private final MongoDbDataContext _dataContext; - private final Map<String, DBCollection> _collections; - private final WriteConcernAdvisor _writeConcernAdvisor; - - public MongoDbUpdateCallback(MongoDbDataContext dataContext, WriteConcernAdvisor writeConcernAdvisor) { - super(dataContext); - _dataContext = dataContext; - _writeConcernAdvisor = writeConcernAdvisor; - _collections = new HashMap<String, DBCollection>(); - } - - @Override - public MongoDbDataContext getDataContext() { - return _dataContext; - } - - public WriteConcernAdvisor getWriteConcernAdvisor() { - return _writeConcernAdvisor; - } - - @Override - public TableCreationBuilder createTable(Schema schema, String name) throws IllegalArgumentException, - IllegalStateException { - return new MongoDbTableCreationBuilder(this, schema, name); - } - - @Override - public RowInsertionBuilder insertInto(Table table) throws IllegalArgumentException, IllegalStateException { - return new MongoDbInsertionBuilder(this, table); - } - - protected void createCollection(String name) { - DBCollection collection = _dataContext.getMongoDb().createCollection(name, new BasicDBObject()); - _collections.put(name, collection); - } - - protected void removeCollection(String name) { - DBCollection collection = getCollection(name); - _collections.remove(name); - collection.drop(); - } - - protected DBCollection getCollection(String name) { - DBCollection collection = _collections.get(name); - if (collection == null) { - collection = _dataContext.getMongoDb().getCollection(name); - _collections.put(name, collection); - } - return collection; - } - - @Override - public void close() { - _collections.clear(); - } - - @Override - public boolean isDropTableSupported() { - return true; - } - - @Override - public TableDropBuilder dropTable(Table table) throws UnsupportedOperationException { - return new MongoDbDropTableBuilder(this, table); - } - - @Override - public boolean isDeleteSupported() { - return true; - } - - @Override - public RowDeletionBuilder deleteFrom(Table table) throws IllegalArgumentException, IllegalStateException, - UnsupportedOperationException { - return new MongoDbDeleteBuilder(this, table); - } -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/SimpleWriteConcernAdvisor.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/SimpleWriteConcernAdvisor.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/SimpleWriteConcernAdvisor.java deleted file mode 100644 index 2d6f76a..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/SimpleWriteConcernAdvisor.java +++ /dev/null @@ -1,50 +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.metamodel.mongodb; - -import com.mongodb.BasicDBObject; -import com.mongodb.DBCollection; -import com.mongodb.WriteConcern; - -/** - * A simple {@link WriteConcernAdvisor} that always returns the same write - * concern. - */ -public class SimpleWriteConcernAdvisor implements WriteConcernAdvisor { - - private final WriteConcern _writeConcern; - - public SimpleWriteConcernAdvisor(WriteConcern writeConcern) { - if (writeConcern == null) { - throw new IllegalArgumentException("WriteConcern cannot be null"); - } - _writeConcern = writeConcern; - } - - @Override - public WriteConcern adviceDeleteQuery(DBCollection collection, BasicDBObject query) { - return _writeConcern; - } - - @Override - public WriteConcern adviceInsert(DBCollection collection, BasicDBObject document) { - return _writeConcern; - } - -} http://git-wip-us.apache.org/repos/asf/metamodel/blob/365fae3a/mongodb/src/main/java/org/apache/metamodel/mongodb/WriteConcernAdvisor.java ---------------------------------------------------------------------- diff --git a/mongodb/src/main/java/org/apache/metamodel/mongodb/WriteConcernAdvisor.java b/mongodb/src/main/java/org/apache/metamodel/mongodb/WriteConcernAdvisor.java deleted file mode 100644 index 006234b..0000000 --- a/mongodb/src/main/java/org/apache/metamodel/mongodb/WriteConcernAdvisor.java +++ /dev/null @@ -1,35 +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.metamodel.mongodb; - -import com.mongodb.BasicDBObject; -import com.mongodb.DBCollection; -import com.mongodb.WriteConcern; - -/** - * Interface for component that advices MetaModel on which {@link WriteConcern} - * to apply to given operations - */ -public interface WriteConcernAdvisor { - - public WriteConcern adviceDeleteQuery(DBCollection collection, BasicDBObject query); - - public WriteConcern adviceInsert(DBCollection collection, BasicDBObject document); - -}
