Repository: nifi Updated Branches: refs/heads/master bff5b7ab7 -> dae2b73d9
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/PutMongoTest.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoTest.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoTest.java deleted file mode 100644 index d0b1a9d..0000000 --- a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/PutMongoTest.java +++ /dev/null @@ -1,311 +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.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.Ignore; -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; - -@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 PutMongoTest 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/RunMongoAggregationIT.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationIT.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationIT.java new file mode 100644 index 0000000..d6c489a --- /dev/null +++ b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationIT.java @@ -0,0 +1,184 @@ +/* + * 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.mongodb.MongoClient; +import com.mongodb.MongoClientURI; +import com.mongodb.client.MongoCollection; +import org.apache.nifi.util.MockFlowFile; +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.io.IOException; +import java.util.Calendar; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +public class RunMongoAggregationIT { + + private static final String MONGO_URI = "mongodb://localhost"; + private static final String DB_NAME = String.format("agg_test-%s", Calendar.getInstance().getTimeInMillis()); + private static final String COLLECTION_NAME = "agg_test_data"; + private static final String AGG_ATTR = "mongo.aggregation.query"; + + private TestRunner runner; + private MongoClient mongoClient; + private Map<String, Integer> mappings; + + @Before + public void setup() { + runner = TestRunners.newTestRunner(RunMongoAggregation.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(RunMongoAggregation.QUERY_ATTRIBUTE, AGG_ATTR); + runner.setValidateExpressionUsage(true); + + mongoClient = new MongoClient(new MongoClientURI(MONGO_URI)); + + MongoCollection<Document> collection = mongoClient.getDatabase(DB_NAME).getCollection(COLLECTION_NAME); + String[] values = new String[] { "a", "b", "c" }; + mappings = new HashMap<>(); + + for (int x = 0; x < values.length; x++) { + for (int y = 0; y < x + 2; y++) { + Document doc = new Document().append("val", values[x]); + collection.insertOne(doc); + } + mappings.put(values[x], x + 2); + } + } + + @After + public void teardown() { + runner = null; + + mongoClient.getDatabase(DB_NAME).drop(); + } + + @Test + public void testAggregation() throws Exception { + final String queryInput = "[\n" + + " {\n" + + " \"$project\": {\n" + + " \"_id\": 0,\n" + + " \"val\": 1\n" + + " }\n" + + " },\n" + + " {\n" + + " \"$group\": {\n" + + " \"_id\": \"$val\",\n" + + " \"doc_count\": {\n" + + " \"$sum\": 1\n" + + " }\n" + + " }\n" + + " }\n" + + "]"; + runner.setProperty(RunMongoAggregation.QUERY, queryInput); + runner.enqueue("test"); + runner.run(1, true, true); + + evaluateRunner(1); + + runner.clearTransferState(); + + runner.setIncomingConnection(false); + runner.run(); //Null parent flowfile + evaluateRunner(0); + + runner.run(); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(RunMongoAggregation.REL_RESULTS); + for (MockFlowFile mff : flowFiles) { + String val = mff.getAttribute(AGG_ATTR); + Assert.assertNotNull("Missing query attribute", val); + Assert.assertEquals("Value was wrong", val, queryInput); + } + } + + @Test + public void testExpressionLanguageSupport() throws Exception { + runner.setVariable("fieldName", "$val"); + runner.setProperty(RunMongoAggregation.QUERY, "[\n" + + " {\n" + + " \"$project\": {\n" + + " \"_id\": 0,\n" + + " \"val\": 1\n" + + " }\n" + + " },\n" + + " {\n" + + " \"$group\": {\n" + + " \"_id\": \"${fieldName}\",\n" + + " \"doc_count\": {\n" + + " \"$sum\": 1\n" + + " }\n" + + " }\n" + + " }\n" + + "]"); + + runner.enqueue("test"); + runner.run(1, true, true); + evaluateRunner(1); + } + + @Test + public void testInvalidQuery(){ + runner.setProperty(RunMongoAggregation.QUERY, "[\n" + + " {\n" + + " \"$invalid_stage\": {\n" + + " \"_id\": 0,\n" + + " \"val\": 1\n" + + " }\n" + + " }\n" + + "]" + ); + runner.enqueue("test"); + runner.run(1, true, true); + runner.assertTransferCount(RunMongoAggregation.REL_RESULTS, 0); + runner.assertTransferCount(RunMongoAggregation.REL_ORIGINAL, 0); + runner.assertTransferCount(RunMongoAggregation.REL_FAILURE, 1); + } + + private void evaluateRunner(int original) throws IOException { + runner.assertTransferCount(RunMongoAggregation.REL_RESULTS, mappings.size()); + runner.assertTransferCount(RunMongoAggregation.REL_ORIGINAL, original); + List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(RunMongoAggregation.REL_RESULTS); + ObjectMapper mapper = new ObjectMapper(); + for (MockFlowFile mockFlowFile : flowFiles) { + byte[] raw = runner.getContentAsByteArray(mockFlowFile); + Map read = mapper.readValue(raw, Map.class); + Assert.assertTrue("Value was not found", mappings.containsKey(read.get("_id"))); + + String queryAttr = mockFlowFile.getAttribute(AGG_ATTR); + Assert.assertNotNull("Query attribute was null.", queryAttr); + Assert.assertTrue("Missing $project", queryAttr.contains("$project")); + Assert.assertTrue("Missing $group", queryAttr.contains("$group")); + } + } +} 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/RunMongoAggregationTest.java ---------------------------------------------------------------------- diff --git a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationTest.java b/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationTest.java deleted file mode 100644 index 95b2200..0000000 --- a/nifi-nar-bundles/nifi-mongodb-bundle/nifi-mongodb-processors/src/test/java/org/apache/nifi/processors/mongodb/RunMongoAggregationTest.java +++ /dev/null @@ -1,186 +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.mongodb.MongoClient; -import com.mongodb.MongoClientURI; -import com.mongodb.client.MongoCollection; -import org.apache.nifi.util.MockFlowFile; -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.io.IOException; -import java.util.Calendar; -import java.util.HashMap; -import java.util.List; -import java.util.Map; - -@Ignore("This is an integration test that requires Mongo to be running.") -public class RunMongoAggregationTest { - - private static final String MONGO_URI = "mongodb://localhost"; - private static final String DB_NAME = String.format("agg_test-%s", Calendar.getInstance().getTimeInMillis()); - private static final String COLLECTION_NAME = "agg_test_data"; - private static final String AGG_ATTR = "mongo.aggregation.query"; - - private TestRunner runner; - private MongoClient mongoClient; - private Map<String, Integer> mappings; - - @Before - public void setup() { - runner = TestRunners.newTestRunner(RunMongoAggregation.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(RunMongoAggregation.QUERY_ATTRIBUTE, AGG_ATTR); - runner.setValidateExpressionUsage(true); - - mongoClient = new MongoClient(new MongoClientURI(MONGO_URI)); - - MongoCollection<Document> collection = mongoClient.getDatabase(DB_NAME).getCollection(COLLECTION_NAME); - String[] values = new String[] { "a", "b", "c" }; - mappings = new HashMap<>(); - - for (int x = 0; x < values.length; x++) { - for (int y = 0; y < x + 2; y++) { - Document doc = new Document().append("val", values[x]); - collection.insertOne(doc); - } - mappings.put(values[x], x + 2); - } - } - - @After - public void teardown() { - runner = null; - - mongoClient.getDatabase(DB_NAME).drop(); - } - - @Test - public void testAggregation() throws Exception { - final String queryInput = "[\n" + - " {\n" + - " \"$project\": {\n" + - " \"_id\": 0,\n" + - " \"val\": 1\n" + - " }\n" + - " },\n" + - " {\n" + - " \"$group\": {\n" + - " \"_id\": \"$val\",\n" + - " \"doc_count\": {\n" + - " \"$sum\": 1\n" + - " }\n" + - " }\n" + - " }\n" + - "]"; - runner.setProperty(RunMongoAggregation.QUERY, queryInput); - runner.enqueue("test"); - runner.run(1, true, true); - - evaluateRunner(1); - - runner.clearTransferState(); - - runner.setIncomingConnection(false); - runner.run(); //Null parent flowfile - evaluateRunner(0); - - runner.run(); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(RunMongoAggregation.REL_RESULTS); - for (MockFlowFile mff : flowFiles) { - String val = mff.getAttribute(AGG_ATTR); - Assert.assertNotNull("Missing query attribute", val); - Assert.assertEquals("Value was wrong", val, queryInput); - } - } - - @Test - public void testExpressionLanguageSupport() throws Exception { - runner.setVariable("fieldName", "$val"); - runner.setProperty(RunMongoAggregation.QUERY, "[\n" + - " {\n" + - " \"$project\": {\n" + - " \"_id\": 0,\n" + - " \"val\": 1\n" + - " }\n" + - " },\n" + - " {\n" + - " \"$group\": {\n" + - " \"_id\": \"${fieldName}\",\n" + - " \"doc_count\": {\n" + - " \"$sum\": 1\n" + - " }\n" + - " }\n" + - " }\n" + - "]"); - - runner.enqueue("test"); - runner.run(1, true, true); - evaluateRunner(1); - } - - @Test - public void testInvalidQuery(){ - runner.setProperty(RunMongoAggregation.QUERY, "[\n" + - " {\n" + - " \"$invalid_stage\": {\n" + - " \"_id\": 0,\n" + - " \"val\": 1\n" + - " }\n" + - " }\n" + - "]" - ); - runner.enqueue("test"); - runner.run(1, true, true); - runner.assertTransferCount(RunMongoAggregation.REL_RESULTS, 0); - runner.assertTransferCount(RunMongoAggregation.REL_ORIGINAL, 0); - runner.assertTransferCount(RunMongoAggregation.REL_FAILURE, 1); - } - - private void evaluateRunner(int original) throws IOException { - runner.assertTransferCount(RunMongoAggregation.REL_RESULTS, mappings.size()); - runner.assertTransferCount(RunMongoAggregation.REL_ORIGINAL, original); - List<MockFlowFile> flowFiles = runner.getFlowFilesForRelationship(RunMongoAggregation.REL_RESULTS); - ObjectMapper mapper = new ObjectMapper(); - for (MockFlowFile mockFlowFile : flowFiles) { - byte[] raw = runner.getContentAsByteArray(mockFlowFile); - Map read = mapper.readValue(raw, Map.class); - Assert.assertTrue("Value was not found", mappings.containsKey(read.get("_id"))); - - String queryAttr = mockFlowFile.getAttribute(AGG_ATTR); - Assert.assertNotNull("Query attribute was null.", queryAttr); - Assert.assertTrue("Missing $project", queryAttr.contains("$project")); - Assert.assertTrue("Missing $group", queryAttr.contains("$group")); - } - } -}
