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

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

commit 8ead1fc58f15805c1d9a556d1db112a00e17e519
Author: Piotr P. Karwasz <[email protected]>
AuthorDate: Wed Sep 2 20:35:56 2026 +0200

    Validate MongoDB sink configuration
---
 .../org/apache/flume/sink/mongodb/MongoDbSink.java | 93 +++++++++++++++++-----
 .../apache/flume/sink/mongodb/TestMongoDbSink.java | 69 ++++++++++++++--
 2 files changed, 134 insertions(+), 28 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..0fcc639 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,22 @@
  */
 package org.apache.flume.sink.mongodb;
 
+import static org.apache.flume.sink.mongodb.MongoDbSinkConstants.BATCH_SIZE;
+import static org.apache.flume.sink.mongodb.MongoDbSinkConstants.COLLECTION;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.COLLECTION_HEADER;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.COLLECTION_MAP_FALLBACK;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.COLLECTION_MAP_PREFIX;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.CONNECTION_URI;
+import static org.apache.flume.sink.mongodb.MongoDbSinkConstants.DATABASE_NAME;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.DEFAULT_BATCH_SIZE;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.DEFAULT_COLLECTION_MAP_FALLBACK;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.DEFAULT_INCLUDE_HEADERS;
+import static 
org.apache.flume.sink.mongodb.MongoDbSinkConstants.INCLUDE_HEADERS;
+import static org.apache.flume.sink.mongodb.MongoDbSinkConstants.WRITE_CONCERN;
+
+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 +54,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 +87,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 +230,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,37 +256,32 @@ 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(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(DATABASE_NAME, 
connectionString.getDatabase());
+        String databaseNameSource = context.containsKey(DATABASE_NAME) ? 
DATABASE_NAME : 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(COLLECTION, 
connectionString.getCollection());
+        String collectionNameSource = context.containsKey(COLLECTION) ? 
COLLECTION : CONNECTION_URI;
+        checkCollectionNameValidity(defaultCollection, collectionNameSource);
 
-        collectionHeader = 
context.getString(MongoDbSinkConstants.COLLECTION_HEADER);
-        collectionMap = 
context.getSubProperties(MongoDbSinkConstants.COLLECTION_MAP_PREFIX);
-        collectionMapFallback = context.getBoolean(
-                MongoDbSinkConstants.COLLECTION_MAP_FALLBACK, 
MongoDbSinkConstants.DEFAULT_COLLECTION_MAP_FALLBACK);
+        collectionHeader = context.getString(COLLECTION_HEADER);
+        collectionMap = context.getSubProperties(COLLECTION_MAP_PREFIX);
+        collectionMap.forEach((key, value) -> 
checkCollectionNameValidity(value, COLLECTION_MAP_PREFIX + key));
+        collectionMapFallback = context.getBoolean(COLLECTION_MAP_FALLBACK, 
DEFAULT_COLLECTION_MAP_FALLBACK);
 
         if (collectionHeader != null && logger.isDebugEnabled()) {
             logger.debug(
                     "Using header {} with mappings {} to select target 
collection", collectionHeader, collectionMap);
         }
 
-        includeHeaders =
-                context.getBoolean(MongoDbSinkConstants.INCLUDE_HEADERS, 
MongoDbSinkConstants.DEFAULT_INCLUDE_HEADERS);
+        includeHeaders = context.getBoolean(INCLUDE_HEADERS, 
DEFAULT_INCLUDE_HEADERS);
 
-        batchSize = context.getInteger(MongoDbSinkConstants.BATCH_SIZE, 
MongoDbSinkConstants.DEFAULT_BATCH_SIZE);
+        batchSize = context.getInteger(BATCH_SIZE, DEFAULT_BATCH_SIZE);
 
-        String writeConcernName = 
context.getString(MongoDbSinkConstants.WRITE_CONCERN);
+        String writeConcernName = context.getString(WRITE_CONCERN);
         if (writeConcernName != null && !writeConcernName.isEmpty()) {
             writeConcern = WriteConcern.valueOf(writeConcernName);
             if (writeConcern == null) {
@@ -282,4 +293,42 @@ 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 `" + CONNECTION_URI + "`");
+        }
+        try {
+            return new ConnectionString(connectionUri);
+        } catch (IllegalArgumentException error) {
+            throw new ConfigurationException(
+                    "Invalid MongoDB connection string in `" + CONNECTION_URI 
+ "`: `" + connectionUri + "`", error);
+        }
+    }
+
+    private static void checkDatabaseNameValidity(@Nullable String 
databaseName, String source) {
+        if (databaseName == null) {
+            throw new ConfigurationException("Missing MongoDB database name; 
set `" + DATABASE_NAME
+                    + "` or include it in `" + 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 `" + COLLECTION
+                    + "` or include it in `" + 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..dac052a 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;
 
@@ -110,28 +111,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