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