NIFI-4929 Converted the majority of MongoDB unit tests to integration tests so they can be reliably run with 'mvn -Pintegration-tests integration-test'
Signed-off-by: Pierre Villard <[email protected]> This closes #2508. Project: http://git-wip-us.apache.org/repos/asf/nifi/repo Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/dae2b73d Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/dae2b73d Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/dae2b73d Branch: refs/heads/master Commit: dae2b73d94f51230ff3193fe7baae657ff10ce82 Parents: bff5b7a Author: Mike Thomsen <[email protected]> Authored: Sat Mar 3 10:17:21 2018 -0500 Committer: Pierre Villard <[email protected]> Committed: Thu Mar 8 19:23:44 2018 +0100 ---------------------------------------------------------------------- .../nifi/processors/mongodb/DeleteMongoIT.java | 118 +++++ .../processors/mongodb/DeleteMongoTest.java | 120 ----- .../nifi/processors/mongodb/GetMongoIT.java | 434 ++++++++++++++++++ .../nifi/processors/mongodb/GetMongoTest.java | 436 ------------------- .../nifi/processors/mongodb/PutMongoIT.java | 309 +++++++++++++ .../processors/mongodb/PutMongoRecordIT.java | 186 ++++++++ .../processors/mongodb/PutMongoRecordTest.java | 188 -------- .../nifi/processors/mongodb/PutMongoTest.java | 311 ------------- .../mongodb/RunMongoAggregationIT.java | 184 ++++++++ .../mongodb/RunMongoAggregationTest.java | 186 -------- 10 files changed, 1231 insertions(+), 1241 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoIT.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoIT.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoIT.java new file mode 100644 index 0000000..42880d3 --- /dev/null +++ b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoIT.java @@ -0,0 +1,118 @@ +/* + * 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.nifi.processors.mongodb; + +import org.bson.Document; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.util.HashMap; +import java.util.Map; + +public class DeleteMongoIT extends MongoWriteTestBase { + @Before + public void setup() { + super.setup(DeleteMongo.class); + collection.insertMany(DOCUMENTS); + } + + @After + public void teardown() { + super.teardown(); + } + + private void testOne(String query, Map<String, String> attrs) { + runner.enqueue(query, attrs); + runner.run(1, true); + runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); + runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); + + Assert.assertEquals("Found a document that should have been deleted.", + 0, collection.count(Document.parse(query))); + } + + @Test + public void testDeleteOne() { + String query = "{ \"_id\": \"doc_1\" }"; + runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ONE); + testOne(query, new HashMap<>()); + Map<String, String> attrs = new HashMap<>(); + attrs.put("mongodb.delete.mode", "one"); + runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ATTR); + query = "{ \"_id\": \"doc_2\" }"; + runner.clearTransferState(); + testOne(query, attrs); + } + + private void manyTest(String query, Map<String, String> attrs) { + runner.enqueue(query, attrs); + runner.run(1, true); + runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); + runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); + + Assert.assertEquals("Found a document that should have been deleted.", + 0, collection.count(Document.parse(query))); + Assert.assertEquals("One document should have been left.", + 1, collection.count(Document.parse("{}"))); + } + + @Test + public void testDeleteMany() { + String query = "{\n" + + "\t\"_id\": {\n" + + "\t\t\"$in\": [\"doc_1\", \"doc_2\"]\n" + + "\t}\n" + + "}"; + runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_MANY); + manyTest(query, new HashMap<>()); + + runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ATTR); + Map<String, String> attrs = new HashMap<>(); + attrs.put("mongodb.delete.mode", "many"); + collection.drop(); + collection.insertMany(DOCUMENTS); + runner.clearTransferState(); + manyTest(query, attrs); + } + + @Test + public void testFailOnNoDeleteOptions() { + String query = "{ \"_id\": \"doc_4\"} "; + runner.enqueue(query); + runner.run(1, true); + runner.assertTransferCount(DeleteMongo.REL_FAILURE, 1); + runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 0); + + Assert.assertEquals("A document was deleted", 3, collection.count(Document.parse("{}"))); + + runner.setProperty(DeleteMongo.FAIL_ON_NO_DELETE, DeleteMongo.NO_FAIL); + runner.clearTransferState(); + runner.enqueue(query); + runner.run(1, true, true); + + + runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); + runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); + + Assert.assertEquals("A document was deleted", 3, collection.count(Document.parse("{}"))); + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoTest.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoTest.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoTest.java deleted file mode 100644 index 00cd55e..0000000 --- a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/DeleteMongoTest.java +++ /dev/null @@ -1,120 +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.nifi.processors.mongodb; - -import org.bson.Document; -import org.junit.After; -import org.junit.Assert; -import org.junit.Before; -import org.junit.Ignore; -import org.junit.Test; - -import java.util.HashMap; -import java.util.Map; - -@Ignore("This is an integration test and should be marked Ignore until someone needs to run it.") -public class DeleteMongoTest extends MongoWriteTestBase { - @Before - public void setup() { - super.setup(DeleteMongo.class); - collection.insertMany(DOCUMENTS); - } - - @After - public void teardown() { - super.teardown(); - } - - private void testOne(String query, Map<String, String> attrs) { - runner.enqueue(query, attrs); - runner.run(1, true); - runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); - runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); - - Assert.assertEquals("Found a document that should have been deleted.", - 0, collection.count(Document.parse(query))); - } - - @Test - public void testDeleteOne() { - String query = "{ \"_id\": \"doc_1\" }"; - runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ONE); - testOne(query, new HashMap<>()); - Map<String, String> attrs = new HashMap<>(); - attrs.put("mongodb.delete.mode", "one"); - runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ATTR); - query = "{ \"_id\": \"doc_2\" }"; - runner.clearTransferState(); - testOne(query, attrs); - } - - private void manyTest(String query, Map<String, String> attrs) { - runner.enqueue(query, attrs); - runner.run(1, true); - runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); - runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); - - Assert.assertEquals("Found a document that should have been deleted.", - 0, collection.count(Document.parse(query))); - Assert.assertEquals("One document should have been left.", - 1, collection.count(Document.parse("{}"))); - } - - @Test - public void testDeleteMany() { - String query = "{\n" + - "\t\"_id\": {\n" + - "\t\t\"$in\": [\"doc_1\", \"doc_2\"]\n" + - "\t}\n" + - "}"; - runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_MANY); - manyTest(query, new HashMap<>()); - - runner.setProperty(DeleteMongo.DELETE_MODE, DeleteMongo.DELETE_ATTR); - Map<String, String> attrs = new HashMap<>(); - attrs.put("mongodb.delete.mode", "many"); - collection.drop(); - collection.insertMany(DOCUMENTS); - runner.clearTransferState(); - manyTest(query, attrs); - } - - @Test - public void testFailOnNoDeleteOptions() { - String query = "{ \"_id\": \"doc_4\"} "; - runner.enqueue(query); - runner.run(1, true); - runner.assertTransferCount(DeleteMongo.REL_FAILURE, 1); - runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 0); - - Assert.assertEquals("A document was deleted", 3, collection.count(Document.parse("{}"))); - - runner.setProperty(DeleteMongo.FAIL_ON_NO_DELETE, DeleteMongo.NO_FAIL); - runner.clearTransferState(); - runner.enqueue(query); - runner.run(1, true, true); - - - runner.assertTransferCount(DeleteMongo.REL_FAILURE, 0); - runner.assertTransferCount(DeleteMongo.REL_SUCCESS, 1); - - Assert.assertEquals("A document was deleted", 3, collection.count(Document.parse("{}"))); - } -} http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoIT.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoIT.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoIT.java new file mode 100644 index 0000000..a079a2b --- /dev/null +++ b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoIT.java @@ -0,0 +1,434 @@ +/* + * 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.nifi.processors.mongodb; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.Lists; +import com.mongodb.MongoClient; +import com.mongodb.MongoClientURI; +import com.mongodb.client.MongoCollection; +import org.apache.nifi.components.ValidationResult; +import org.apache.nifi.flowfile.attributes.CoreAttributes; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.util.MockFlowFile; +import org.apache.nifi.util.MockProcessContext; +import org.apache.nifi.util.TestRunner; +import org.apache.nifi.util.TestRunners; +import org.bson.Document; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.text.SimpleDateFormat; +import java.util.Calendar; +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; +import java.util.Map; + +public class GetMongoIT { + private static final String MONGO_URI = "mongodb://localhost"; + private static final String DB_NAME = GetMongoIT.class.getSimpleName().toLowerCase(); + private static final String COLLECTION_NAME = "test"; + + private static final List<Document> DOCUMENTS; + private static final Calendar CAL; + + static { + CAL = Calendar.getInstance(); + DOCUMENTS = Lists.newArrayList( + new Document("_id", "doc_1").append("a", 1).append("b", 2).append("c", 3), + new Document("_id", "doc_2").append("a", 1).append("b", 2).append("c", 4).append("date_field", CAL.getTime()), + new Document("_id", "doc_3").append("a", 1).append("b", 3) + ); + } + + private TestRunner runner; + private MongoClient mongoClient; + + @Before + public void setup() { + runner = TestRunners.newTestRunner(GetMongo.class); + runner.setVariable("uri", MONGO_URI); + runner.setVariable("db", DB_NAME); + runner.setVariable("collection", COLLECTION_NAME); + runner.setProperty(AbstractMongoProcessor.URI, "${uri}"); + runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, "${db}"); + runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, "${collection}"); + runner.setProperty(GetMongo.USE_PRETTY_PRINTING, GetMongo.YES_PP); + runner.setIncomingConnection(false); + + mongoClient = new MongoClient(new MongoClientURI(MONGO_URI)); + + MongoCollection<Document> collection = mongoClient.getDatabase(DB_NAME).getCollection(COLLECTION_NAME); + collection.insertMany(DOCUMENTS); + } + + @After + public void teardown() { + runner = null; + + mongoClient.getDatabase(DB_NAME).drop(); + } + + @Test + public void testValidators() { + + TestRunner runner = TestRunners.newTestRunner(GetMongo.class); + Collection<ValidationResult> results; + ProcessContext pc; + + // missing uri, db, collection + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(3, results.size()); + Iterator<ValidationResult> it = results.iterator(); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo URI is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Database Name is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Collection Name is required")); + + // missing query - is ok + runner.setProperty(AbstractMongoProcessor.URI, MONGO_URI); + runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, DB_NAME); + runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, COLLECTION_NAME); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(0, results.size()); + + // invalid query + runner.setProperty(GetMongo.QUERY, "{a: x,y,z}"); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(1, results.size()); + Assert.assertTrue(results.iterator().next().toString().contains("is invalid because")); + + // invalid projection + runner.setVariable("projection", "{a: x,y,z}"); + runner.setProperty(GetMongo.QUERY, "{a: 1}"); + runner.setProperty(GetMongo.PROJECTION, "{a: z}"); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(1, results.size()); + Assert.assertTrue(results.iterator().next().toString().contains("is invalid")); + + // invalid sort + runner.removeProperty(GetMongo.PROJECTION); + runner.setProperty(GetMongo.SORT, "{a: x,y,z}"); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(1, results.size()); + Assert.assertTrue(results.iterator().next().toString().contains("is invalid")); + } + + @Test + public void testCleanJson() throws Exception { + runner.setVariable("query", "{\"_id\": \"doc_2\"}"); + runner.setProperty(GetMongo.QUERY, "${query}"); + runner.setProperty(GetMongo.JSON_TYPE, GetMongo.JSON_STANDARD); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + byte[] raw = runner.getContentAsByteArray(flowFiles.get(0)); + ObjectMapper mapper = new ObjectMapper(); + Map<String, Object> parsed = mapper.readValue(raw, Map.class); + SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd"); + + Assert.assertTrue(parsed.get("date_field").getClass() == String.class); + Assert.assertTrue(((String)parsed.get("date_field")).startsWith(format.format(CAL.getTime()))); + } + + @Test + public void testReadOneDocument() throws Exception { + runner.setVariable("query", "{a: 1, b: 3}"); + runner.setProperty(GetMongo.QUERY, "${query}"); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + flowFiles.get(0).assertContentEquals(DOCUMENTS.get(2).toJson()); + } + + @Test + public void testReadMultipleDocuments() throws Exception { + runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 3); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + for (int i=0; i < flowFiles.size(); i++) { + flowFiles.get(i).assertContentEquals(DOCUMENTS.get(i).toJson()); + } + } + + @Test + public void testProjection() throws Exception { + runner.setProperty(GetMongo.QUERY, "{a: 1, b: 3}"); + runner.setProperty(GetMongo.PROJECTION, "{_id: 0, a: 1}"); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + Document expected = new Document("a", 1); + flowFiles.get(0).assertContentEquals(expected.toJson()); + } + + @Test + public void testSort() throws Exception { + runner.setVariable("sort", "{a: -1, b: -1, c: 1}"); + runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); + runner.setProperty(GetMongo.SORT, "${sort}"); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 3); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + flowFiles.get(0).assertContentEquals(DOCUMENTS.get(2).toJson()); + flowFiles.get(1).assertContentEquals(DOCUMENTS.get(0).toJson()); + flowFiles.get(2).assertContentEquals(DOCUMENTS.get(1).toJson()); + } + + @Test + public void testLimit() throws Exception { + runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); + runner.setProperty(GetMongo.LIMIT, "${limit}"); + runner.setVariable("limit", "1"); + runner.run(); + + runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + flowFiles.get(0).assertContentEquals(DOCUMENTS.get(0).toJson()); + } + + @Test + public void testResultsPerFlowfile() throws Exception { + runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "${results.per.flowfile}"); + runner.setVariable("results.per.flowfile", "2"); + runner.enqueue("{}"); + runner.setIncomingConnection(true); + runner.run(); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 2); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + List<MockFlowFile> results = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + Assert.assertTrue("Flowfile was empty", results.get(0).getSize() > 0); + Assert.assertEquals("Wrong mime type", results.get(0).getAttribute(CoreAttributes.MIME_TYPE.key()), "application/json"); + } + + @Test + public void testBatchSize() throws Exception { + runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "2"); + runner.setProperty(GetMongo.BATCH_SIZE, "${batch.size}"); + runner.setVariable("batch.size", "1"); + runner.enqueue("{}"); + runner.setIncomingConnection(true); + runner.run(); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 2); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + List<MockFlowFile> results = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + Assert.assertTrue("Flowfile was empty", results.get(0).getSize() > 0); + Assert.assertEquals("Wrong mime type", results.get(0).getAttribute(CoreAttributes.MIME_TYPE.key()), "application/json"); + } + + @Test + public void testConfigurablePrettyPrint() { + runner.setProperty(GetMongo.JSON_TYPE, GetMongo.JSON_STANDARD); + runner.setProperty(GetMongo.LIMIT, "1"); + runner.enqueue("{}"); + runner.setIncomingConnection(true); + runner.run(); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + byte[] raw = runner.getContentAsByteArray(flowFiles.get(0)); + String json = new String(raw); + Assert.assertTrue("JSON did not have new lines.", json.contains("\n")); + runner.clearTransferState(); + runner.setProperty(GetMongo.USE_PRETTY_PRINTING, GetMongo.NO_PP); + runner.enqueue("{}"); + runner.run(); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + raw = runner.getContentAsByteArray(flowFiles.get(0)); + json = new String(raw); + Assert.assertFalse("New lines detected", json.contains("\n")); + } + + private void testQueryAttribute(String attr, String expected) { + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); + for (MockFlowFile mff : flowFiles) { + String val = mff.getAttribute(attr); + Assert.assertNotNull("Missing query attribute", val); + Assert.assertEquals("Value was wrong", expected, val); + } + } + + @Test + public void testQueryAttribute() { + /* + * Test original behavior; Manually set query of {}, no input + */ + final String attr = "query.attr"; + runner.setProperty(GetMongo.QUERY, "{}"); + runner.setProperty(GetMongo.QUERY_ATTRIBUTE, attr); + runner.run(); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 3); + testQueryAttribute(attr, "{}"); + + runner.clearTransferState(); + + /* + * Test original behavior; No Input/Empty val = {} + */ + runner.removeProperty(GetMongo.QUERY); + runner.setIncomingConnection(false); + runner.run(); + testQueryAttribute(attr, "{}"); + + runner.clearTransferState(); + + /* + * Input flowfile with {} as the query + */ + + runner.setIncomingConnection(true); + runner.enqueue("{}"); + runner.run(); + testQueryAttribute(attr, "{}"); + + /* + * Input flowfile with invalid query + */ + + runner.clearTransferState(); + runner.enqueue("invalid query"); + runner.run(); + + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); + runner.assertTransferCount(GetMongo.REL_FAILURE, 1); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 0); + } + + /* + * Query read behavior tests + */ + @Test + public void testReadQueryFromBodyWithEL() { + Map attributes = new HashMap(); + attributes.put("field", "c"); + attributes.put("value", "4"); + String query = "{ \"${field}\": { \"$gte\": ${value}}}"; + runner.setIncomingConnection(true); + runner.setProperty(GetMongo.QUERY, query); + runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "10"); + runner.setValidateExpressionUsage(true); + runner.enqueue("test", attributes); + runner.run(1, true, true); + + runner.assertTransferCount(GetMongo.REL_FAILURE, 0); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + } + + @Test + public void testReadQueryFromBodyNoEL() { + String query = "{ \"c\": { \"$gte\": 4 }}"; + runner.setIncomingConnection(true); + runner.removeProperty(GetMongo.QUERY); + runner.enqueue(query); + runner.run(1, true, true); + + runner.assertTransferCount(GetMongo.REL_FAILURE, 0); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + + } + + @Test + public void testReadQueryFromQueryParamNoConnection() { + String query = "{ \"c\": { \"$gte\": 4 }}"; + runner.setProperty(GetMongo.QUERY, query); + runner.setIncomingConnection(false); + runner.run(1, true, true); + runner.assertTransferCount(GetMongo.REL_FAILURE, 0); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + + } + + @Test + public void testReadQueryFromQueryParamWithConnection() { + String query = "{ \"c\": { \"$gte\": ${value} }}"; + Map<String, String> attrs = new HashMap<>(); + attrs.put("value", "4"); + + runner.setProperty(GetMongo.QUERY, query); + runner.setIncomingConnection(true); + runner.enqueue("test", attrs); + runner.run(1, true, true); + runner.assertTransferCount(GetMongo.REL_FAILURE, 0); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); + + } + + @Test + public void testQueryParamMissingWithNoFlowfile() { + Exception ex = null; + + try { + runner.assertValid(); + runner.setIncomingConnection(false); + runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "1"); + runner.run(1, true, true); + } catch (Exception pe) { + ex = pe; + } + + Assert.assertNull("An exception was thrown!", ex); + runner.assertTransferCount(GetMongo.REL_FAILURE, 0); + runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); + runner.assertTransferCount(GetMongo.REL_SUCCESS, 3); + } + /* + * End query read behavior tests + */ +} http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoTest.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoTest.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoTest.java deleted file mode 100644 index 6cf8d62..0000000 --- a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/GetMongoTest.java +++ /dev/null @@ -1,436 +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.nifi.processors.mongodb; - -import com.fasterxml.jackson.databind.ObjectMapper; -import com.google.common.collect.Lists; -import com.mongodb.MongoClient; -import com.mongodb.MongoClientURI; -import com.mongodb.client.MongoCollection; -import org.apache.nifi.components.ValidationResult; -import org.apache.nifi.flowfile.attributes.CoreAttributes; -import org.apache.nifi.processor.ProcessContext; -import org.apache.nifi.util.MockFlowFile; -import org.apache.nifi.util.MockProcessContext; -import org.apache.nifi.util.TestRunner; -import org.apache.nifi.util.TestRunners; -import org.bson.Document; -import org.junit.After; -import org.junit.Assert; -import org.junit.Before; -import org.junit.Ignore; -import org.junit.Test; - -import java.text.SimpleDateFormat; -import java.util.Calendar; -import java.util.Collection; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Iterator; -import java.util.List; -import java.util.Map; - -@Ignore("Integration tests that cause failures in some environments. Require that they be run from Maven to run the embedded mongo maven plugin. Maven Plugin also fails in my CentOS 7 environment.") -public class GetMongoTest { - private static final String MONGO_URI = "mongodb://localhost"; - private static final String DB_NAME = GetMongoTest.class.getSimpleName().toLowerCase(); - private static final String COLLECTION_NAME = "test"; - - private static final List<Document> DOCUMENTS; - private static final Calendar CAL; - - static { - CAL = Calendar.getInstance(); - DOCUMENTS = Lists.newArrayList( - new Document("_id", "doc_1").append("a", 1).append("b", 2).append("c", 3), - new Document("_id", "doc_2").append("a", 1).append("b", 2).append("c", 4).append("date_field", CAL.getTime()), - new Document("_id", "doc_3").append("a", 1).append("b", 3) - ); - } - - private TestRunner runner; - private MongoClient mongoClient; - - @Before - public void setup() { - runner = TestRunners.newTestRunner(GetMongo.class); - runner.setVariable("uri", MONGO_URI); - runner.setVariable("db", DB_NAME); - runner.setVariable("collection", COLLECTION_NAME); - runner.setProperty(AbstractMongoProcessor.URI, "${uri}"); - runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, "${db}"); - runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, "${collection}"); - runner.setProperty(GetMongo.USE_PRETTY_PRINTING, GetMongo.YES_PP); - runner.setIncomingConnection(false); - - mongoClient = new MongoClient(new MongoClientURI(MONGO_URI)); - - MongoCollection<Document> collection = mongoClient.getDatabase(DB_NAME).getCollection(COLLECTION_NAME); - collection.insertMany(DOCUMENTS); - } - - @After - public void teardown() { - runner = null; - - mongoClient.getDatabase(DB_NAME).drop(); - } - - @Test - public void testValidators() { - - TestRunner runner = TestRunners.newTestRunner(GetMongo.class); - Collection<ValidationResult> results; - ProcessContext pc; - - // missing uri, db, collection - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(3, results.size()); - Iterator<ValidationResult> it = results.iterator(); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo URI is required")); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Database Name is required")); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Collection Name is required")); - - // missing query - is ok - runner.setProperty(AbstractMongoProcessor.URI, MONGO_URI); - runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, DB_NAME); - runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, COLLECTION_NAME); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(0, results.size()); - - // invalid query - runner.setProperty(GetMongo.QUERY, "{a: x,y,z}"); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(1, results.size()); - Assert.assertTrue(results.iterator().next().toString().contains("is invalid because")); - - // invalid projection - runner.setVariable("projection", "{a: x,y,z}"); - runner.setProperty(GetMongo.QUERY, "{a: 1}"); - runner.setProperty(GetMongo.PROJECTION, "{a: z}"); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(1, results.size()); - Assert.assertTrue(results.iterator().next().toString().contains("is invalid")); - - // invalid sort - runner.removeProperty(GetMongo.PROJECTION); - runner.setProperty(GetMongo.SORT, "{a: x,y,z}"); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(1, results.size()); - Assert.assertTrue(results.iterator().next().toString().contains("is invalid")); - } - - @Test - public void testCleanJson() throws Exception { - runner.setVariable("query", "{\"_id\": \"doc_2\"}"); - runner.setProperty(GetMongo.QUERY, "${query}"); - runner.setProperty(GetMongo.JSON_TYPE, GetMongo.JSON_STANDARD); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - byte[] raw = runner.getContentAsByteArray(flowFiles.get(0)); - ObjectMapper mapper = new ObjectMapper(); - Map<String, Object> parsed = mapper.readValue(raw, Map.class); - SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd"); - - Assert.assertTrue(parsed.get("date_field").getClass() == String.class); - Assert.assertTrue(((String)parsed.get("date_field")).startsWith(format.format(CAL.getTime()))); - } - - @Test - public void testReadOneDocument() throws Exception { - runner.setVariable("query", "{a: 1, b: 3}"); - runner.setProperty(GetMongo.QUERY, "${query}"); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - flowFiles.get(0).assertContentEquals(DOCUMENTS.get(2).toJson()); - } - - @Test - public void testReadMultipleDocuments() throws Exception { - runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 3); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - for (int i=0; i < flowFiles.size(); i++) { - flowFiles.get(i).assertContentEquals(DOCUMENTS.get(i).toJson()); - } - } - - @Test - public void testProjection() throws Exception { - runner.setProperty(GetMongo.QUERY, "{a: 1, b: 3}"); - runner.setProperty(GetMongo.PROJECTION, "{_id: 0, a: 1}"); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - Document expected = new Document("a", 1); - flowFiles.get(0).assertContentEquals(expected.toJson()); - } - - @Test - public void testSort() throws Exception { - runner.setVariable("sort", "{a: -1, b: -1, c: 1}"); - runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); - runner.setProperty(GetMongo.SORT, "${sort}"); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 3); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - flowFiles.get(0).assertContentEquals(DOCUMENTS.get(2).toJson()); - flowFiles.get(1).assertContentEquals(DOCUMENTS.get(0).toJson()); - flowFiles.get(2).assertContentEquals(DOCUMENTS.get(1).toJson()); - } - - @Test - public void testLimit() throws Exception { - runner.setProperty(GetMongo.QUERY, "{a: {$exists: true}}"); - runner.setProperty(GetMongo.LIMIT, "${limit}"); - runner.setVariable("limit", "1"); - runner.run(); - - runner.assertAllFlowFilesTransferred(GetMongo.REL_SUCCESS, 1); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - flowFiles.get(0).assertContentEquals(DOCUMENTS.get(0).toJson()); - } - - @Test - public void testResultsPerFlowfile() throws Exception { - runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "${results.per.flowfile}"); - runner.setVariable("results.per.flowfile", "2"); - runner.enqueue("{}"); - runner.setIncomingConnection(true); - runner.run(); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 2); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - List<MockFlowFile> results = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - Assert.assertTrue("Flowfile was empty", results.get(0).getSize() > 0); - Assert.assertEquals("Wrong mime type", results.get(0).getAttribute(CoreAttributes.MIME_TYPE.key()), "application/json"); - } - - @Test - public void testBatchSize() throws Exception { - runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "2"); - runner.setProperty(GetMongo.BATCH_SIZE, "${batch.size}"); - runner.setVariable("batch.size", "1"); - runner.enqueue("{}"); - runner.setIncomingConnection(true); - runner.run(); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 2); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - List<MockFlowFile> results = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - Assert.assertTrue("Flowfile was empty", results.get(0).getSize() > 0); - Assert.assertEquals("Wrong mime type", results.get(0).getAttribute(CoreAttributes.MIME_TYPE.key()), "application/json"); - } - - @Test - public void testConfigurablePrettyPrint() { - runner.setProperty(GetMongo.JSON_TYPE, GetMongo.JSON_STANDARD); - runner.setProperty(GetMongo.LIMIT, "1"); - runner.enqueue("{}"); - runner.setIncomingConnection(true); - runner.run(); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - byte[] raw = runner.getContentAsByteArray(flowFiles.get(0)); - String json = new String(raw); - Assert.assertTrue("JSON did not have new lines.", json.contains("\n")); - runner.clearTransferState(); - runner.setProperty(GetMongo.USE_PRETTY_PRINTING, GetMongo.NO_PP); - runner.enqueue("{}"); - runner.run(); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - raw = runner.getContentAsByteArray(flowFiles.get(0)); - json = new String(raw); - Assert.assertFalse("New lines detected", json.contains("\n")); - } - - private void testQueryAttribute(String attr, String expected) { - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(GetMongo.REL_SUCCESS); - for (MockFlowFile mff : flowFiles) { - String val = mff.getAttribute(attr); - Assert.assertNotNull("Missing query attribute", val); - Assert.assertEquals("Value was wrong", expected, val); - } - } - - @Test - public void testQueryAttribute() { - /* - * Test original behavior; Manually set query of {}, no input - */ - final String attr = "query.attr"; - runner.setProperty(GetMongo.QUERY, "{}"); - runner.setProperty(GetMongo.QUERY_ATTRIBUTE, attr); - runner.run(); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 3); - testQueryAttribute(attr, "{}"); - - runner.clearTransferState(); - - /* - * Test original behavior; No Input/Empty val = {} - */ - runner.removeProperty(GetMongo.QUERY); - runner.setIncomingConnection(false); - runner.run(); - testQueryAttribute(attr, "{}"); - - runner.clearTransferState(); - - /* - * Input flowfile with {} as the query - */ - - runner.setIncomingConnection(true); - runner.enqueue("{}"); - runner.run(); - testQueryAttribute(attr, "{}"); - - /* - * Input flowfile with invalid query - */ - - runner.clearTransferState(); - runner.enqueue("invalid query"); - runner.run(); - - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); - runner.assertTransferCount(GetMongo.REL_FAILURE, 1); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 0); - } - - /* - * Query read behavior tests - */ - @Test - public void testReadQueryFromBodyWithEL() { - Map attributes = new HashMap(); - attributes.put("field", "c"); - attributes.put("value", "4"); - String query = "{ \"${field}\": { \"$gte\": ${value}}}"; - runner.setIncomingConnection(true); - runner.setProperty(GetMongo.QUERY, query); - runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "10"); - runner.setValidateExpressionUsage(true); - runner.enqueue("test", attributes); - runner.run(1, true, true); - - runner.assertTransferCount(GetMongo.REL_FAILURE, 0); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - } - - @Test - public void testReadQueryFromBodyNoEL() { - String query = "{ \"c\": { \"$gte\": 4 }}"; - runner.setIncomingConnection(true); - runner.removeProperty(GetMongo.QUERY); - runner.enqueue(query); - runner.run(1, true, true); - - runner.assertTransferCount(GetMongo.REL_FAILURE, 0); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - - } - - @Test - public void testReadQueryFromQueryParamNoConnection() { - String query = "{ \"c\": { \"$gte\": 4 }}"; - runner.setProperty(GetMongo.QUERY, query); - runner.setIncomingConnection(false); - runner.run(1, true, true); - runner.assertTransferCount(GetMongo.REL_FAILURE, 0); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - - } - - @Test - public void testReadQueryFromQueryParamWithConnection() { - String query = "{ \"c\": { \"$gte\": ${value} }}"; - Map<String, String> attrs = new HashMap<>(); - attrs.put("value", "4"); - - runner.setProperty(GetMongo.QUERY, query); - runner.setIncomingConnection(true); - runner.enqueue("test", attrs); - runner.run(1, true, true); - runner.assertTransferCount(GetMongo.REL_FAILURE, 0); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 1); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 1); - - } - - @Test - public void testQueryParamMissingWithNoFlowfile() { - Exception ex = null; - - try { - runner.assertValid(); - runner.setIncomingConnection(false); - runner.setProperty(GetMongo.RESULTS_PER_FLOWFILE, "1"); - runner.run(1, true, true); - } catch (Exception pe) { - ex = pe; - } - - Assert.assertNull("An exception was thrown!", ex); - runner.assertTransferCount(GetMongo.REL_FAILURE, 0); - runner.assertTransferCount(GetMongo.REL_ORIGINAL, 0); - runner.assertTransferCount(GetMongo.REL_SUCCESS, 3); - } - /* - * End query read behavior tests - */ -} http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoIT.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoIT.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoIT.java new file mode 100644 index 0000000..7e7f330 --- /dev/null +++ b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoIT.java @@ -0,0 +1,309 @@ +/* + * 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.nifi.processors.mongodb; + +import org.apache.nifi.components.ValidationResult; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.util.MockFlowFile; +import org.apache.nifi.util.MockProcessContext; +import org.apache.nifi.util.TestRunner; +import org.apache.nifi.util.TestRunners; +import org.bson.Document; +import org.bson.types.ObjectId; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.nio.charset.StandardCharsets; +import java.util.Collection; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; + +import static org.junit.Assert.assertEquals; + +public class PutMongoIT extends MongoWriteTestBase { + @Before + public void setup() { + super.setup(PutMongo.class); + } + + @After + public void teardown() { + super.teardown(); + } + + private byte[] documentToByteArray(Document doc) { + return doc.toJson().getBytes(StandardCharsets.UTF_8); + } + + @Test + public void testValidators() { + TestRunner runner = TestRunners.newTestRunner(PutMongo.class); + Collection<ValidationResult> results; + ProcessContext pc; + + // missing uri, db, collection + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(3, results.size()); + Iterator<ValidationResult> it = results.iterator(); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo URI is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Database Name is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Collection Name is required")); + + // invalid write concern + runner.setProperty(AbstractMongoProcessor.URI, MONGO_URI); + runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, DATABASE_NAME); + runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, COLLECTION_NAME); + runner.setProperty(PutMongo.WRITE_CONCERN, "xyz"); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(1, results.size()); + Assert.assertTrue(results.iterator().next().toString().matches("'Write Concern' .* is invalid because Given value not found in allowed set .*")); + + // valid write concern + runner.setProperty(PutMongo.WRITE_CONCERN, PutMongo.WRITE_CONCERN_UNACKNOWLEDGED); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(0, results.size()); + } + + @Test + public void testInsertOne() throws Exception { + Document doc = DOCUMENTS.get(0); + byte[] bytes = documentToByteArray(doc); + + runner.enqueue(bytes); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongo.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS).get(0); + out.assertContentEquals(bytes); + + // verify 1 doc inserted into the collection + assertEquals(1, collection.count()); + assertEquals(doc, collection.find().first()); + } + + @Test + public void testInsertMany() throws Exception { + for (Document doc : DOCUMENTS) { + runner.enqueue(documentToByteArray(doc)); + } + runner.run(3); + + runner.assertAllFlowFilesTransferred(PutMongo.REL_SUCCESS, 3); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS); + for (int i=0; i < flowFiles.size(); i++) { + flowFiles.get(i).assertContentEquals(DOCUMENTS.get(i).toJson()); + } + + // verify 3 docs inserted into the collection + assertEquals(3, collection.count()); + } + + @Test + public void testInsertWithDuplicateKey() throws Exception { + // pre-insert one document + collection.insertOne(DOCUMENTS.get(0)); + + for (Document doc : DOCUMENTS) { + runner.enqueue(documentToByteArray(doc)); + } + runner.run(3); + + // first doc failed, other 2 succeeded + runner.assertTransferCount(PutMongo.REL_FAILURE, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongo.REL_FAILURE).get(0); + out.assertContentEquals(documentToByteArray(DOCUMENTS.get(0))); + + runner.assertTransferCount(PutMongo.REL_SUCCESS, 2); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS); + for (int i=0; i < flowFiles.size(); i++) { + flowFiles.get(i).assertContentEquals(DOCUMENTS.get(i+1).toJson()); + } + + // verify 2 docs inserted into the collection for a total of 3 + assertEquals(3, collection.count()); + } + + /** + * Verifies that 'update' does not insert if 'upsert' if false. + * @see #testUpsert() + */ + @Test + public void testUpdateDoesNotInsert() throws Exception { + Document doc = DOCUMENTS.get(0); + byte[] bytes = documentToByteArray(doc); + + runner.setProperty(PutMongo.MODE, "update"); + runner.enqueue(bytes); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongo.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS).get(0); + out.assertContentEquals(bytes); + + // nothing was in collection, so nothing to update since upsert defaults to false + assertEquals(0, collection.count()); + } + + /** + * Verifies that 'update' does insert if 'upsert' is true. + * @see #testUpdateDoesNotInsert() + */ + @Test + public void testUpsert() throws Exception { + Document doc = DOCUMENTS.get(0); + byte[] bytes = documentToByteArray(doc); + + runner.setProperty(PutMongo.MODE, "update"); + runner.setProperty(PutMongo.UPSERT, "true"); + runner.enqueue(bytes); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongo.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS).get(0); + out.assertContentEquals(bytes); + + // verify 1 doc inserted into the collection + assertEquals(1, collection.count()); + assertEquals(doc, collection.find().first()); + } + + @Test + public void testUpdate() throws Exception { + Document doc = DOCUMENTS.get(0); + + // pre-insert document + collection.insertOne(doc); + + // modify the object + doc.put("abc", "123"); + doc.put("xyz", "456"); + doc.remove("c"); + + byte[] bytes = documentToByteArray(doc); + + runner.setProperty(PutMongo.MODE, "update"); + runner.enqueue(bytes); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongo.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongo.REL_SUCCESS).get(0); + out.assertContentEquals(bytes); + + assertEquals(1, collection.count()); + assertEquals(doc, collection.find().first()); + } + + @Test + public void testUpsertWithOperators() throws Exception { + String upsert = "{\n" + + " \"_id\": \"Test\",\n" + + " \"$push\": {\n" + + " \"testArr\": { \"msg\": \"Hi\" }\n" + + " }\n" + + "}"; + runner.setProperty(PutMongo.UPDATE_MODE, PutMongo.UPDATE_WITH_OPERATORS); + runner.setProperty(PutMongo.MODE, "update"); + runner.setProperty(PutMongo.UPSERT, "true"); + for (int x = 0; x < 3; x++) { + runner.enqueue(upsert.getBytes()); + } + runner.run(3, true, true); + runner.assertTransferCount(PutMongo.REL_FAILURE, 0); + runner.assertTransferCount(PutMongo.REL_SUCCESS, 3); + + Document query = new Document("_id", "Test"); + Document result = collection.find(query).first(); + List array = (List)result.get("testArr"); + Assert.assertNotNull("Array was empty", array); + Assert.assertEquals("Wrong size", array.size(), 3); + for (int index = 0; index < array.size(); index++) { + Document doc = (Document)array.get(index); + String msg = doc.getString("msg"); + Assert.assertNotNull("Msg was null", msg); + Assert.assertEquals("Msg had wrong value", msg, "Hi"); + } + } + + /* + * Start NIFI-4759 Regression Tests + * + * 2 issues with ID field: + * + * * Assumed _id is the update key, causing failures when the user configured a different one in the UI. + * * Treated _id as a string even when it is an ObjectID sent from another processor as a string value. + * + * Expected behavior: + * + * * update key field should work no matter what (legal) value it is set to be. + * * _ids that are ObjectID should become real ObjectIDs when added to Mongo. + * * _ids that are arbitrary strings should be still go in as strings. + * + */ + @Test + public void testNiFi_4759_Regressions() { + String[] upserts = new String[]{ + "{ \"_id\": \"12345\", \"$set\": { \"msg\": \"Hello, world\" } }", + "{ \"_id\": \"5a5617b9c1f5de6d8276e87d\", \"$set\": { \"msg\": \"Hello, world\" } }", + "{ \"updateKey\": \"12345\", \"$set\": { \"msg\": \"Hello, world\" } }" + }; + + String[] updateKeyProps = new String[] { "_id", "_id", "updateKey" }; + Object[] updateKeys = new Object[] { "12345", new ObjectId("5a5617b9c1f5de6d8276e87d"), "12345" }; + int index = 0; + + runner.setProperty(PutMongo.UPDATE_MODE, PutMongo.UPDATE_WITH_OPERATORS); + runner.setProperty(PutMongo.MODE, "update"); + runner.setProperty(PutMongo.UPSERT, "true"); + + final int LIMIT = 2; + + for (String upsert : upserts) { + runner.setProperty(PutMongo.UPDATE_QUERY_KEY, updateKeyProps[index]); + for (int x = 0; x < LIMIT; x++) { + runner.enqueue(upsert); + } + runner.run(LIMIT, true, true); + runner.assertTransferCount(PutMongo.REL_FAILURE, 0); + runner.assertTransferCount(PutMongo.REL_SUCCESS, LIMIT); + + Document query = new Document(updateKeyProps[index], updateKeys[index]); + Document result = collection.find(query).first(); + Assert.assertNotNull("Result was null", result); + Assert.assertEquals("Count was wrong", 1, collection.count(query)); + runner.clearTransferState(); + index++; + } + } +} http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordIT.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordIT.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordIT.java new file mode 100644 index 0000000..db4be57 --- /dev/null +++ b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordIT.java @@ -0,0 +1,186 @@ +/* + * 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.nifi.processors.mongodb; + + +import org.apache.nifi.components.ValidationResult; +import org.apache.nifi.processor.ProcessContext; +import org.apache.nifi.serialization.SimpleRecordSchema; +import org.apache.nifi.serialization.record.MapRecord; +import org.apache.nifi.serialization.record.MockRecordParser; +import org.apache.nifi.serialization.record.RecordField; +import org.apache.nifi.serialization.record.RecordFieldType; +import org.apache.nifi.serialization.record.RecordSchema; +import org.apache.nifi.util.MockFlowFile; +import org.apache.nifi.util.MockProcessContext; +import org.apache.nifi.util.TestRunner; +import org.apache.nifi.util.TestRunners; +import org.bson.Document; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; + +import static org.junit.Assert.assertEquals; + +public class PutMongoRecordIT extends MongoWriteTestBase { + + private MockRecordParser recordReader; + + @Before + public void setup() throws Exception { + super.setup(PutMongoRecord.class); + recordReader = new MockRecordParser(); + runner.addControllerService("reader", recordReader); + runner.enableControllerService(recordReader); + runner.setProperty(PutMongoRecord.RECORD_READER_FACTORY, "reader"); + } + + @After + public void teardown() { + super.teardown(); + } + + private byte[] documentToByteArray(Document doc) { + return doc.toJson().getBytes(StandardCharsets.UTF_8); + } + + @Test + public void testValidators() throws Exception { + TestRunner runner = TestRunners.newTestRunner(PutMongoRecord.class); + runner.addControllerService("reader", recordReader); + runner.enableControllerService(recordReader); + Collection<ValidationResult> results; + ProcessContext pc; + + // missing uri, db, collection, RecordReader + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(4, results.size()); + Iterator<ValidationResult> it = results.iterator(); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo URI is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Database Name is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Collection Name is required")); + Assert.assertTrue(it.next().toString().contains("is invalid because Record Reader is required")); + + // invalid write concern + runner.setProperty(AbstractMongoProcessor.URI, MONGO_URI); + runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, DATABASE_NAME); + runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, COLLECTION_NAME); + runner.setProperty(PutMongoRecord.RECORD_READER_FACTORY, "reader"); + runner.setProperty(PutMongoRecord.WRITE_CONCERN, "xyz"); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(1, results.size()); + Assert.assertTrue(results.iterator().next().toString().matches("'Write Concern' .* is invalid because Given value not found in allowed set .*")); + + // valid write concern + runner.setProperty(PutMongoRecord.WRITE_CONCERN, PutMongoRecord.WRITE_CONCERN_UNACKNOWLEDGED); + runner.enqueue(new byte[0]); + pc = runner.getProcessContext(); + results = new HashSet<>(); + if (pc instanceof MockProcessContext) { + results = ((MockProcessContext) pc).validate(); + } + Assert.assertEquals(0, results.size()); + } + + @Test + public void testInsertFlatRecords() throws Exception { + recordReader.addSchemaField("name", RecordFieldType.STRING); + recordReader.addSchemaField("age", RecordFieldType.INT); + recordReader.addSchemaField("sport", RecordFieldType.STRING); + + recordReader.addRecord("John Doe", 48, "Soccer"); + recordReader.addRecord("Jane Doe", 47, "Tennis"); + recordReader.addRecord("Sally Doe", 47, "Curling"); + recordReader.addRecord("Jimmy Doe", 14, null); + recordReader.addRecord("Pizza Doe", 14, null); + + runner.enqueue(""); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongoRecord.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongoRecord.REL_SUCCESS).get(0); + + + // verify 1 doc inserted into the collection + assertEquals(5, collection.count()); + //assertEquals(doc, collection.find().first()); + } + + @Test + public void testInsertNestedRecords() throws Exception { + recordReader.addSchemaField("id", RecordFieldType.INT); + final List<RecordField> personFields = new ArrayList<>(); + final RecordField nameField = new RecordField("name", RecordFieldType.STRING.getDataType()); + final RecordField ageField = new RecordField("age", RecordFieldType.INT.getDataType()); + final RecordField sportField = new RecordField("sport", RecordFieldType.STRING.getDataType()); + personFields.add(nameField); + personFields.add(ageField); + personFields.add(sportField); + final RecordSchema personSchema = new SimpleRecordSchema(personFields); + recordReader.addSchemaField("person", RecordFieldType.RECORD); + recordReader.addRecord(1, new MapRecord(personSchema, new HashMap<String,Object>() {{ + put("name", "John Doe"); + put("age", 48); + put("sport", "Soccer"); + }})); + recordReader.addRecord(2, new MapRecord(personSchema, new HashMap<String,Object>() {{ + put("name", "Jane Doe"); + put("age", 47); + put("sport", "Tennis"); + }})); + recordReader.addRecord(3, new MapRecord(personSchema, new HashMap<String,Object>() {{ + put("name", "Sally Doe"); + put("age", 47); + put("sport", "Curling"); + }})); + recordReader.addRecord(4, new MapRecord(personSchema, new HashMap<String,Object>() {{ + put("name", "Jimmy Doe"); + put("age", 14); + put("sport", null); + }})); + + runner.enqueue(""); + runner.run(); + + runner.assertAllFlowFilesTransferred(PutMongoRecord.REL_SUCCESS, 1); + MockFlowFile out = runner.getFlowFilesForRelationship(PutMongoRecord.REL_SUCCESS).get(0); + + + // verify 1 doc inserted into the collection + assertEquals(4, collection.count()); + //assertEquals(doc, collection.find().first()); + } +} \ No newline at end of file http://git-wip-us.apache.org/repos/asf/nifi/blob/dae2b73d/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordTest.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordTest.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordTest.java deleted file mode 100644 index a8cbf82..0000000 --- a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoRecordTest.java +++ /dev/null @@ -1,188 +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.nifi.processors.mongodb; - - -import org.apache.nifi.components.ValidationResult; -import org.apache.nifi.processor.ProcessContext; -import org.apache.nifi.serialization.SimpleRecordSchema; -import org.apache.nifi.serialization.record.MapRecord; -import org.apache.nifi.serialization.record.MockRecordParser; -import org.apache.nifi.serialization.record.RecordField; -import org.apache.nifi.serialization.record.RecordFieldType; -import org.apache.nifi.serialization.record.RecordSchema; -import org.apache.nifi.util.MockFlowFile; -import org.apache.nifi.util.MockProcessContext; -import org.apache.nifi.util.TestRunner; -import org.apache.nifi.util.TestRunners; -import org.bson.Document; -import org.junit.After; -import org.junit.Assert; -import org.junit.Before; -import org.junit.Ignore; -import org.junit.Test; - -import java.nio.charset.StandardCharsets; -import java.util.ArrayList; -import java.util.Collection; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Iterator; -import java.util.List; - -import static org.junit.Assert.assertEquals; - -@Ignore("Integration tests that cause failures in some environments") -public class PutMongoRecordTest extends MongoWriteTestBase { - - private MockRecordParser recordReader; - - @Before - public void setup() throws Exception { - super.setup(PutMongoRecord.class); - recordReader = new MockRecordParser(); - runner.addControllerService("reader", recordReader); - runner.enableControllerService(recordReader); - runner.setProperty(PutMongoRecord.RECORD_READER_FACTORY, "reader"); - } - - @After - public void teardown() { - super.teardown(); - } - - private byte[] documentToByteArray(Document doc) { - return doc.toJson().getBytes(StandardCharsets.UTF_8); - } - - @Test - public void testValidators() throws Exception { - TestRunner runner = TestRunners.newTestRunner(PutMongoRecord.class); - runner.addControllerService("reader", recordReader); - runner.enableControllerService(recordReader); - Collection<ValidationResult> results; - ProcessContext pc; - - // missing uri, db, collection, RecordReader - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(4, results.size()); - Iterator<ValidationResult> it = results.iterator(); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo URI is required")); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Database Name is required")); - Assert.assertTrue(it.next().toString().contains("is invalid because Mongo Collection Name is required")); - Assert.assertTrue(it.next().toString().contains("is invalid because Record Reader is required")); - - // invalid write concern - runner.setProperty(AbstractMongoProcessor.URI, MONGO_URI); - runner.setProperty(AbstractMongoProcessor.DATABASE_NAME, DATABASE_NAME); - runner.setProperty(AbstractMongoProcessor.COLLECTION_NAME, COLLECTION_NAME); - runner.setProperty(PutMongoRecord.RECORD_READER_FACTORY, "reader"); - runner.setProperty(PutMongoRecord.WRITE_CONCERN, "xyz"); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(1, results.size()); - Assert.assertTrue(results.iterator().next().toString().matches("'Write Concern' .* is invalid because Given value not found in allowed set .*")); - - // valid write concern - runner.setProperty(PutMongoRecord.WRITE_CONCERN, PutMongoRecord.WRITE_CONCERN_UNACKNOWLEDGED); - runner.enqueue(new byte[0]); - pc = runner.getProcessContext(); - results = new HashSet<>(); - if (pc instanceof MockProcessContext) { - results = ((MockProcessContext) pc).validate(); - } - Assert.assertEquals(0, results.size()); - } - - @Test - public void testInsertFlatRecords() throws Exception { - recordReader.addSchemaField("name", RecordFieldType.STRING); - recordReader.addSchemaField("age", RecordFieldType.INT); - recordReader.addSchemaField("sport", RecordFieldType.STRING); - - recordReader.addRecord("John Doe", 48, "Soccer"); - recordReader.addRecord("Jane Doe", 47, "Tennis"); - recordReader.addRecord("Sally Doe", 47, "Curling"); - recordReader.addRecord("Jimmy Doe", 14, null); - recordReader.addRecord("Pizza Doe", 14, null); - - runner.enqueue(""); - runner.run(); - - runner.assertAllFlowFilesTransferred(PutMongoRecord.REL_SUCCESS, 1); - MockFlowFile out = runner.getFlowFilesForRelationship(PutMongoRecord.REL_SUCCESS).get(0); - - - // verify 1 doc inserted into the collection - assertEquals(5, collection.count()); - //assertEquals(doc, collection.find().first()); - } - - @Test - public void testInsertNestedRecords() throws Exception { - recordReader.addSchemaField("id", RecordFieldType.INT); - final List<RecordField> personFields = new ArrayList<>(); - final RecordField nameField = new RecordField("name", RecordFieldType.STRING.getDataType()); - final RecordField ageField = new RecordField("age", RecordFieldType.INT.getDataType()); - final RecordField sportField = new RecordField("sport", RecordFieldType.STRING.getDataType()); - personFields.add(nameField); - personFields.add(ageField); - personFields.add(sportField); - final RecordSchema personSchema = new SimpleRecordSchema(personFields); - recordReader.addSchemaField("person", RecordFieldType.RECORD); - recordReader.addRecord(1, new MapRecord(personSchema, new HashMap<String,Object>() {{ - put("name", "John Doe"); - put("age", 48); - put("sport", "Soccer"); - }})); - recordReader.addRecord(2, new MapRecord(personSchema, new HashMap<String,Object>() {{ - put("name", "Jane Doe"); - put("age", 47); - put("sport", "Tennis"); - }})); - recordReader.addRecord(3, new MapRecord(personSchema, new HashMap<String,Object>() {{ - put("name", "Sally Doe"); - put("age", 47); - put("sport", "Curling"); - }})); - recordReader.addRecord(4, new MapRecord(personSchema, new HashMap<String,Object>() {{ - put("name", "Jimmy Doe"); - put("age", 14); - put("sport", null); - }})); - - runner.enqueue(""); - runner.run(); - - runner.assertAllFlowFilesTransferred(PutMongoRecord.REL_SUCCESS, 1); - MockFlowFile out = runner.getFlowFilesForRelationship(PutMongoRecord.REL_SUCCESS).get(0); - - - // verify 1 doc inserted into the collection - assertEquals(4, collection.count()); - //assertEquals(doc, collection.find().first()); - } -} \ No newline at end of file
