This is an automated email from the ASF dual-hosted git repository.

ppkarwasz pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/logging-flume-mongodb.git


The following commit(s) were added to refs/heads/main by this push:
     new 03c8bdb  Validate MongoDB sink configuration (#10)
03c8bdb is described below

commit 03c8bdb7602d0a8e3c5a059b5e9b223aa5ef16e6
Author: Piotr P. Karwasz <[email protected]>
AuthorDate: Thu Sep 3 07:47:13 2026 +0200

    Validate MongoDB sink configuration (#10)
    
    Validate MongoDB sink configuration using MongoDB client API.
---
 .../org/apache/flume/sink/mongodb/MongoDbSink.java | 75 ++++++++++++++++++----
 .../apache/flume/sink/mongodb/TestMongoDbSink.java | 74 ++++++++++++++++++---
 2 files changed, 125 insertions(+), 24 deletions(-)

diff --git 
a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSink.java
 
b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSink.java
index 9c32289..518916e 100644
--- 
a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSink.java
+++ 
b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSink.java
@@ -16,7 +16,9 @@
  */
 package org.apache.flume.sink.mongodb;
 
+import com.mongodb.ConnectionString;
 import com.mongodb.MongoException;
+import com.mongodb.MongoNamespace;
 import com.mongodb.WriteConcern;
 import com.mongodb.client.MongoClient;
 import com.mongodb.client.MongoClients;
@@ -39,6 +41,7 @@ import org.apache.flume.sink.AbstractSink;
 import org.apache.logging.log4j.LogManager;
 import org.apache.logging.log4j.Logger;
 import org.bson.Document;
+import org.jspecify.annotations.Nullable;
 
 /**
  * A Flume Sink that writes events to MongoDB.
@@ -71,7 +74,7 @@ public class MongoDbSink extends AbstractSink implements 
Configurable, BatchSize
 
     private static final Logger logger = 
LogManager.getLogger(MongoDbSink.class);
 
-    private String connectionUri;
+    private ConnectionString connectionString;
     private String databaseName;
     private String defaultCollection;
     private String collectionHeader;
@@ -214,7 +217,7 @@ public class MongoDbSink extends AbstractSink implements 
Configurable, BatchSize
 
     @Override
     public synchronized void start() {
-        mongoClient = MongoClients.create(connectionUri);
+        mongoClient = MongoClients.create(connectionString);
         MongoDatabase mongoDatabase = mongoClient.getDatabase(databaseName);
         writer = new DefaultMongoDbWriter(mongoDatabase, writeConcern);
         counter.start();
@@ -240,23 +243,25 @@ public class MongoDbSink extends AbstractSink implements 
Configurable, BatchSize
 
     @Override
     public void configure(Context context) {
-        connectionUri = context.getString(MongoDbSinkConstants.CONNECTION_URI);
-        if (connectionUri == null || connectionUri.isEmpty()) {
-            throw new ConfigurationException("mongodb.uri must be specified");
-        }
+        String connectionUri = 
context.getString(MongoDbSinkConstants.CONNECTION_URI);
+        connectionString = createConnectionString(connectionUri);
 
-        databaseName = context.getString(MongoDbSinkConstants.DATABASE_NAME);
-        if (databaseName == null || databaseName.isEmpty()) {
-            throw new ConfigurationException("mongodb.database must be 
specified");
-        }
+        databaseName = context.getString(MongoDbSinkConstants.DATABASE_NAME, 
connectionString.getDatabase());
+        String databaseNameSource = 
context.containsKey(MongoDbSinkConstants.DATABASE_NAME)
+                ? MongoDbSinkConstants.DATABASE_NAME
+                : MongoDbSinkConstants.CONNECTION_URI;
+        checkDatabaseNameValidity(databaseName, databaseNameSource);
 
-        defaultCollection = context.getString(MongoDbSinkConstants.COLLECTION);
-        if (defaultCollection == null || defaultCollection.isEmpty()) {
-            throw new ConfigurationException("mongodb.collection must be 
specified");
-        }
+        defaultCollection = context.getString(MongoDbSinkConstants.COLLECTION, 
connectionString.getCollection());
+        String collectionNameSource = 
context.containsKey(MongoDbSinkConstants.COLLECTION)
+                ? MongoDbSinkConstants.COLLECTION
+                : MongoDbSinkConstants.CONNECTION_URI;
+        checkCollectionNameValidity(defaultCollection, collectionNameSource);
 
         collectionHeader = 
context.getString(MongoDbSinkConstants.COLLECTION_HEADER);
         collectionMap = 
context.getSubProperties(MongoDbSinkConstants.COLLECTION_MAP_PREFIX);
+        collectionMap.forEach(
+                (key, value) -> checkCollectionNameValidity(value, 
MongoDbSinkConstants.COLLECTION_MAP_PREFIX + key));
         collectionMapFallback = context.getBoolean(
                 MongoDbSinkConstants.COLLECTION_MAP_FALLBACK, 
MongoDbSinkConstants.DEFAULT_COLLECTION_MAP_FALLBACK);
 
@@ -282,4 +287,46 @@ public class MongoDbSink extends AbstractSink implements 
Configurable, BatchSize
             counter = new SinkCounter(getName());
         }
     }
+
+    private static ConnectionString createConnectionString(@Nullable String 
connectionUri) {
+        if (connectionUri == null) {
+            throw new ConfigurationException(
+                    "Missing MongoDB connection string in `" + 
MongoDbSinkConstants.CONNECTION_URI + "`");
+        }
+        try {
+            return new ConnectionString(connectionUri);
+        } catch (IllegalArgumentException error) {
+            throw new ConfigurationException(
+                    "Invalid MongoDB connection string in `" + 
MongoDbSinkConstants.CONNECTION_URI + "`: `"
+                            + connectionUri + "`",
+                    error);
+        }
+    }
+
+    private static void checkDatabaseNameValidity(@Nullable String 
databaseName, String source) {
+        if (databaseName == null) {
+            throw new ConfigurationException("Missing MongoDB database name; 
set `" + MongoDbSinkConstants.DATABASE_NAME
+                    + "` or include it in `" + 
MongoDbSinkConstants.CONNECTION_URI + "`");
+        }
+        try {
+            MongoNamespace.checkDatabaseNameValidity(databaseName);
+        } catch (IllegalArgumentException error) {
+            throw new ConfigurationException(
+                    "Invalid MongoDB database name in `" + source + "`: `" + 
databaseName + "`", error);
+        }
+    }
+
+    private static void checkCollectionNameValidity(@Nullable String 
collectionName, String source) {
+        if (collectionName == null) {
+            throw new ConfigurationException("Missing MongoDB collection name; 
set `"
+                    + MongoDbSinkConstants.COLLECTION + "` or include it in `" 
+ MongoDbSinkConstants.CONNECTION_URI
+                    + "`");
+        }
+        try {
+            MongoNamespace.checkCollectionNameValidity(collectionName);
+        } catch (IllegalArgumentException error) {
+            throw new ConfigurationException(
+                    "Invalid MongoDB collection name in `" + source + "`: `" + 
collectionName + "`", error);
+        }
+    }
 }
diff --git 
a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSink.java
 
b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSink.java
index 3dc4de0..c1fa57e 100644
--- 
a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSink.java
+++ 
b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSink.java
@@ -18,6 +18,7 @@ package org.apache.flume.sink.mongodb;
 
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
 
@@ -49,7 +50,6 @@ public class TestMongoDbSink {
      */
     private static final class FakeMongoDbWriter implements MongoDbWriter {
         private final Map<String, List<Document>> written = new 
LinkedHashMap<>();
-        private boolean closed = false;
 
         @Override
         public void write(String collectionName, List<Document> documents) {
@@ -57,9 +57,7 @@ public class TestMongoDbSink {
         }
 
         @Override
-        public void close() {
-            closed = true;
-        }
+        public void close() {}
     }
 
     private static Context baseContext() {
@@ -110,28 +108,84 @@ public class TestMongoDbSink {
         tx.close();
     }
 
-    @Test(expected = ConfigurationException.class)
+    private static void assertConfigurationFailure(Context context, String 
expectedMessage) {
+        ConfigurationException error =
+                assertThrows(ConfigurationException.class, () -> new 
MongoDbSink().configure(context));
+        assertEquals(expectedMessage, error.getMessage());
+    }
+
+    @Test
+    public void testConfigureDatabaseAndCollectionFromUri() {
+        Context context = new Context();
+        context.put(MongoDbSinkConstants.CONNECTION_URI, 
"mongodb://localhost:27017/uriDatabase.uriCollection");
+
+        MongoDbSink sink = new MongoDbSink();
+        sink.configure(context);
+
+        assertEquals("uriDatabase", sink.getDatabaseName());
+        assertEquals("uriCollection", sink.getDefaultCollection());
+        assertEquals(MongoDbSinkConstants.DEFAULT_BATCH_SIZE, 
sink.getBatchSize());
+    }
+
+    @Test
+    public void testConfigureExplicitNamesOverrideUri() {
+        Context context = new Context();
+        context.put(MongoDbSinkConstants.CONNECTION_URI, 
"mongodb://localhost:27017/uriDatabase.uriCollection");
+        context.put(MongoDbSinkConstants.DATABASE_NAME, "configuredDatabase");
+        context.put(MongoDbSinkConstants.COLLECTION, "configuredCollection");
+
+        MongoDbSink sink = new MongoDbSink();
+        sink.configure(context);
+
+        assertEquals("configuredDatabase", sink.getDatabaseName());
+        assertEquals("configuredCollection", sink.getDefaultCollection());
+    }
+
+    @Test
     public void testConfigureMissingUri() {
         Context context = new Context();
         context.put(MongoDbSinkConstants.DATABASE_NAME, "testDb");
         context.put(MongoDbSinkConstants.COLLECTION, "col");
-        new MongoDbSink().configure(context);
+        assertConfigurationFailure(context, "Missing MongoDB connection string 
in `mongodb.uri`");
     }
 
-    @Test(expected = ConfigurationException.class)
+    @Test
     public void testConfigureMissingDatabase() {
         Context context = new Context();
         context.put(MongoDbSinkConstants.CONNECTION_URI, 
"mongodb://localhost:27017");
         context.put(MongoDbSinkConstants.COLLECTION, "col");
-        new MongoDbSink().configure(context);
+        assertConfigurationFailure(
+                context, "Missing MongoDB database name; set 
`mongodb.database` or include it in `mongodb.uri`");
     }
 
-    @Test(expected = ConfigurationException.class)
+    @Test
     public void testConfigureMissingCollection() {
         Context context = new Context();
         context.put(MongoDbSinkConstants.CONNECTION_URI, 
"mongodb://localhost:27017");
         context.put(MongoDbSinkConstants.DATABASE_NAME, "testDb");
-        new MongoDbSink().configure(context);
+        assertConfigurationFailure(
+                context, "Missing MongoDB collection name; set 
`mongodb.collection` or include it in `mongodb.uri`");
+    }
+
+    @Test
+    public void testConfigureInvalidUri() {
+        Context context = baseContext();
+        context.put(MongoDbSinkConstants.CONNECTION_URI, "invalid");
+        assertConfigurationFailure(context, "Invalid MongoDB connection string 
in `mongodb.uri`: `invalid`");
+    }
+
+    @Test
+    public void testConfigureInvalidDatabase() {
+        Context context = baseContext();
+        context.put(MongoDbSinkConstants.DATABASE_NAME, "invalid/database");
+        assertConfigurationFailure(context, "Invalid MongoDB database name in 
`mongodb.database`: `invalid/database`");
+    }
+
+    @Test
+    public void testConfigureEmptyMappedCollection() {
+        Context context = baseContext();
+        context.put(MongoDbSinkConstants.COLLECTION_MAP_PREFIX + "typeA", "");
+        assertConfigurationFailure(context, "Invalid MongoDB collection name 
in `mongodb.collectionMap.typeA`: ``");
     }
 
     @Test

Reply via email to