From 09d6432822165187b768148537e10d6e94dc9643 Mon Sep 17 00:00:00 2001 From: Ralph Goers Date: Tue, 1 Sep 2026 21:14:10 -0700 Subject: [PATCH 1/6] handle duplicate events --- .flattened-pom.xml | 216 ++++++++++++++++++ flume-mongodb-sink/pom.xml | 7 + .../sink/mongodb/DefaultMongoDbWriter.java | 38 ++- .../flume/sink/mongodb/MongoDbSink.java | 23 +- .../sink/mongodb/MongoDbSinkCounter.java | 50 ++++ .../sink/mongodb/MongoDbWriteResult.java | 45 ++++ .../flume/sink/mongodb/MongoDbWriter.java | 10 +- .../flume/sink/mongodb/TestMongoDbSink.java | 26 ++- .../sink/mongodb/TestMongoDbSinkEmbedded.java | 164 +++++++++++++ 9 files changed, 566 insertions(+), 13 deletions(-) create mode 100644 .flattened-pom.xml create mode 100644 flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSinkCounter.java create mode 100644 flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriteResult.java create mode 100644 flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java diff --git a/.flattened-pom.xml b/.flattened-pom.xml new file mode 100644 index 0000000..815ad8e --- /dev/null +++ b/.flattened-pom.xml @@ -0,0 +1,216 @@ + + + + 4.0.0 + + org.apache.flume + flume-parent + 2.0.0-SNAPSHOT + + + org.apache.flume + flume-mongodb-parent + 2.0.0-SNAPSHOT + pom + Flume MongoDB Parent + https://logging.apache.org/flume/2.x/index.html/flume-mongodb-parent + 2022 + + Apache Software Foundation + http://www.apache.org + + + + The Apache Software License, Version 2.0 + http://www.apache.org/licenses/LICENSE-2.0.txt + + + + + rgoers + Ralph Goers + rgoers@apache.org + Intuit + + + + + Flume User List + user-subscribe@flume.apache.org + user-unsubscribe@flume.apache.org + user@flume.apache.org + http://mail-archives.apache.org/mod_mbox/flume-user/ + + + Flume Developer List + dev-subscribe@flume.apache.org + dev-unsubscribe@flume.apache.org + dev@flume.apache.org + http://mail-archives.apache.org/mod_mbox/flume-dev/ + + + Flume Commits + commits-subscribe@flume.apache.org + commits-unsubscribe@flume.apache.org + commits@flume.apache.org + http://mail-archives.apache.org/mod_mbox/flume-commits/ + + + + flume-mongodb-sink + + + https://gitbox.apache.org/repos/asf/flume-spring-boot.git + https://gitbox.apache.org/repos/asf/flume-spring-boot.git + https://gitbox.apache.org/repos/asf/flume-spring-boot.git + + + JIRA + https://issues.apache.org/jira/browse/FLUME + + + 5.10.0 + 2.17.0 + 1.12.0 + 2.0.0 + Ralph Goers + 0.12 + 2.26.1 + 4.13.2 + B3D8E1BA + 11 + 2.27.2 + 1.6 + 1.11 + UTF-8 + org.apache.flume.mongodb + 2.9 + 1.9.0 + 4.1.18 + 2.0.0-SNAPSHOT + 4.7.2.1 + 11 + rgoers@apache.org + + + + + org.apache.flume + flume-ng-core + ${flume.version} + + + org.apache.flume + flume-ng-sdk + ${flume.version} + + + org.apache.flume + flume-ng-configuration + ${flume.version} + + + org.mongodb + mongodb-driver-bom + ${mongodb.version} + pom + import + + + io.dropwizard.metrics + metrics-core + ${dropwizard-metrics.version} + + + org.apache.logging.log4j + log4j-api + ${log4j.version} + + + org.apache.logging.log4j + log4j-core + ${log4j.version} + + + com.fasterxml.jackson.core + jackson-core + ${jackson.version} + + + junit + junit + ${junit.version} + + + org.mockito + mockito-all + ${mockito.version} + test + + + + + + + org.apache.rat + apache-rat-plugin + ${rat.version} + + + verify.rat + verify + + check + + + + + + **/.idea/ + **/*.iml + src/main/resources/META-INF/services/**/* + **/nb-configuration.xml + .git/ + patchprocess/ + .gitignore + **/*.yml + **/*.yaml + **/*.json + .repository/ + **/*.diff + **/*.patch + **/*.avsc + **/*.avro + **/docs/** + **/test/resources/** + **/.settings/* + **/.classpath + **/.project + **/target/** + **/derby.log + **/metastore_db/ + .mvn/** + **/exclude-pmd.properties + + true + + + + + diff --git a/flume-mongodb-sink/pom.xml b/flume-mongodb-sink/pom.xml index ddcb4f9..b221f7e 100644 --- a/flume-mongodb-sink/pom.xml +++ b/flume-mongodb-sink/pom.xml @@ -76,6 +76,13 @@ test + + de.flapdoodle.embed + de.flapdoodle.embed.mongo + 3.5.4 + test + + diff --git a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java index d039717..28fb4ab 100644 --- a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java +++ b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java @@ -16,10 +16,15 @@ */ package org.apache.flume.sink.mongodb; +import com.mongodb.DuplicateKeyException; +import com.mongodb.ErrorCategory; +import com.mongodb.MongoWriteException; import com.mongodb.WriteConcern; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; import java.util.List; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; import org.bson.Document; /** @@ -28,6 +33,8 @@ */ public class DefaultMongoDbWriter implements MongoDbWriter { + private static final Logger logger = LogManager.getLogger(DefaultMongoDbWriter.class); + private final MongoDatabase mongoDatabase; private final WriteConcern writeConcern; @@ -37,12 +44,39 @@ public DefaultMongoDbWriter(MongoDatabase mongoDatabase, WriteConcern writeConce } @Override - public void write(String collectionName, List documents) { + public MongoDbWriteResult write(String collectionName, List documents) { MongoCollection collection = mongoDatabase.getCollection(collectionName); if (writeConcern != null) { collection = collection.withWriteConcern(writeConcern); } - collection.insertMany(documents); + + long insertedCount = 0; + long duplicateCount = 0; + // Insert one document at a time (rather than insertMany) so that a + // single duplicate key does not abort the rest of the batch. + for (Document document : documents) { + try { + collection.insertOne(document); + insertedCount++; + } catch (DuplicateKeyException ex) { + logger.warn( + "Duplicate key while inserting into collection {}, skipping event: {}", + collectionName, + ex.getMessage()); + duplicateCount++; + } catch (MongoWriteException ex) { + if (ex.getError().getCategory() == ErrorCategory.DUPLICATE_KEY) { + logger.warn( + "Duplicate key while inserting into collection {}, skipping event: {}", + collectionName, + ex.getMessage()); + duplicateCount++; + } else { + throw ex; + } + } + } + return new MongoDbWriteResult(insertedCount, duplicateCount); } @Override 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..adc545a 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 @@ -34,7 +34,6 @@ import org.apache.flume.conf.BatchSizeSupported; import org.apache.flume.conf.Configurable; import org.apache.flume.conf.ConfigurationException; -import org.apache.flume.instrumentation.SinkCounter; import org.apache.flume.sink.AbstractSink; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -83,7 +82,7 @@ public class MongoDbSink extends AbstractSink implements Configurable, BatchSize private MongoClient mongoClient; private MongoDbWriter writer; - private SinkCounter counter; + private MongoDbSinkCounter counter; // For testing public String getDatabaseName() { @@ -94,6 +93,11 @@ public String getDefaultCollection() { return defaultCollection; } + /** For testing: number of events skipped due to duplicate keys. */ + public long getDuplicateEventCount() { + return counter.getDuplicateEventCount(); + } + @Override public long getBatchSize() { return batchSize; @@ -138,12 +142,19 @@ public Status process() throws EventDeliveryException { .add(document); } + long insertedEvents = 0; + long duplicateEvents = 0; for (Map.Entry> entry : documentsByCollection.entrySet()) { - writer.write(entry.getKey(), entry.getValue()); + MongoDbWriteResult writeResult = writer.write(entry.getKey(), entry.getValue()); + insertedEvents += writeResult.getInsertedCount(); + duplicateEvents += writeResult.getDuplicateCount(); } - if (processedEvents > 0) { - counter.addToEventDrainSuccessCount(processedEvents); + if (insertedEvents > 0) { + counter.addToEventDrainSuccessCount(insertedEvents); + } + if (duplicateEvents > 0) { + counter.addToDuplicateEventCount(duplicateEvents); } transaction.commit(); @@ -279,7 +290,7 @@ public void configure(Context context) { } if (counter == null) { - counter = new SinkCounter(getName()); + counter = new MongoDbSinkCounter(getName()); } } } diff --git a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSinkCounter.java b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSinkCounter.java new file mode 100644 index 0000000..97737d2 --- /dev/null +++ b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbSinkCounter.java @@ -0,0 +1,50 @@ +/* + * 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.flume.sink.mongodb; + +import org.apache.flume.instrumentation.SinkCounter; + +/** + * {@link SinkCounter} extension that additionally tracks the number of + * events skipped because they duplicated a document already present in + * MongoDB (i.e. resulted in a {@code com.mongodb.DuplicateKeyException}). + * These events are not counted towards {@code eventDrainSuccessCount} since + * they were not actually inserted, but they should also not be treated as a + * batch failure. + */ +public class MongoDbSinkCounter extends SinkCounter { + + private static final String COUNTER_DUPLICATE_EVENT = "sink.event.duplicate"; + + private static final String[] ATTRIBUTES = {COUNTER_DUPLICATE_EVENT}; + + public MongoDbSinkCounter(String name) { + super(name, ATTRIBUTES); + } + + public long getDuplicateEventCount() { + return get(COUNTER_DUPLICATE_EVENT); + } + + public long incrementDuplicateEventCount() { + return increment(COUNTER_DUPLICATE_EVENT); + } + + public long addToDuplicateEventCount(long delta) { + return addAndGet(COUNTER_DUPLICATE_EVENT, delta); + } +} diff --git a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriteResult.java b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriteResult.java new file mode 100644 index 0000000..d697119 --- /dev/null +++ b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriteResult.java @@ -0,0 +1,45 @@ +/* + * 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.flume.sink.mongodb; + +/** + * Outcome of writing a batch of documents to a single MongoDB collection. + */ +public final class MongoDbWriteResult { + + private final long insertedCount; + private final long duplicateCount; + + public MongoDbWriteResult(long insertedCount, long duplicateCount) { + this.insertedCount = insertedCount; + this.duplicateCount = duplicateCount; + } + + /** Number of documents that were successfully inserted. */ + public long getInsertedCount() { + return insertedCount; + } + + /** + * Number of documents that were skipped because they duplicated a + * document that already existed in the collection (i.e. triggered a + * {@code com.mongodb.DuplicateKeyException}). + */ + public long getDuplicateCount() { + return duplicateCount; + } +} diff --git a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriter.java b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriter.java index 005fb53..d458bc6 100644 --- a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriter.java +++ b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/MongoDbWriter.java @@ -28,12 +28,18 @@ public interface MongoDbWriter { /** - * Writes the given documents to the named collection. + * Writes the given documents to the named collection. Documents that + * fail to insert because they duplicate an existing document (i.e. + * trigger a {@code com.mongodb.DuplicateKeyException}) are skipped + * rather than causing the whole batch to fail; they are reported via + * {@link MongoDbWriteResult#getDuplicateCount()}. * * @param collectionName the target collection name * @param documents the documents to insert, in order + * @return the number of documents inserted and the number skipped as + * duplicates */ - void write(String collectionName, List documents); + MongoDbWriteResult write(String collectionName, List documents); /** * Releases any resources (e.g. the underlying MongoDB client) held by 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..a9a65c2 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 @@ -37,7 +37,6 @@ import org.apache.flume.conf.Configurables; import org.apache.flume.conf.ConfigurationException; import org.apache.flume.event.EventBuilder; -import org.apache.flume.instrumentation.SinkCounter; import org.bson.Document; import org.junit.Test; @@ -50,10 +49,14 @@ public class TestMongoDbSink { private static final class FakeMongoDbWriter implements MongoDbWriter { private final Map> written = new LinkedHashMap<>(); private boolean closed = false; + private long duplicateCountToReport = 0; @Override - public void write(String collectionName, List documents) { + public MongoDbWriteResult write(String collectionName, List documents) { written.computeIfAbsent(collectionName, k -> new ArrayList<>()).addAll(documents); + long duplicates = Math.min(duplicateCountToReport, documents.size()); + duplicateCountToReport -= duplicates; + return new MongoDbWriteResult(documents.size() - duplicates, duplicates); } @Override @@ -98,7 +101,7 @@ private static MongoDbSink createSink(Context context, MongoDbWriter writer) { channel.start(); Configurables.configure(sink, context); setInternalState(sink, "writer", writer); - setInternalState(sink, "counter", new SinkCounter("test")); + setInternalState(sink, "counter", new MongoDbSinkCounter("test")); return sink; } @@ -134,6 +137,23 @@ public void testConfigureMissingCollection() { new MongoDbSink().configure(context); } + @Test + public void testDuplicateEventsDoNotFailBatchAndAreCountedSeparately() throws EventDeliveryException { + FakeMongoDbWriter writer = new FakeMongoDbWriter(); + writer.duplicateCountToReport = 1; + Context context = baseContext(); + MongoDbSink sink = createSink(context, writer); + Channel channel = sink.getChannel(); + + putEvent(channel, "{\"foo\":\"1\"}".getBytes(StandardCharsets.UTF_8), new HashMap()); + putEvent(channel, "{\"foo\":\"2\"}".getBytes(StandardCharsets.UTF_8), new HashMap()); + + Sink.Status status = sink.process(); + + assertEquals(Sink.Status.READY, status); + assertEquals(1, sink.getDuplicateEventCount()); + } + @Test public void testWritesToDefaultCollectionWhenNoHeaderConfigured() throws EventDeliveryException { FakeMongoDbWriter writer = new FakeMongoDbWriter(); diff --git a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java new file mode 100644 index 0000000..cc4ffbf --- /dev/null +++ b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java @@ -0,0 +1,164 @@ +/* + * 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.flume.sink.mongodb; + +import static org.junit.Assert.assertEquals; + +import com.mongodb.client.MongoClient; +import com.mongodb.client.MongoClients; +import com.mongodb.client.MongoCollection; +import com.mongodb.client.model.IndexOptions; +import com.mongodb.client.model.Indexes; +import de.flapdoodle.embed.mongo.MongodExecutable; +import de.flapdoodle.embed.mongo.MongodProcess; +import de.flapdoodle.embed.mongo.MongodStarter; +import de.flapdoodle.embed.mongo.config.MongodConfig; +import de.flapdoodle.embed.mongo.config.Net; +import de.flapdoodle.embed.mongo.distribution.Version; +import de.flapdoodle.embed.process.runtime.Network; +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import org.apache.flume.Channel; +import org.apache.flume.Context; +import org.apache.flume.EventDeliveryException; +import org.apache.flume.Sink; +import org.apache.flume.Transaction; +import org.apache.flume.channel.MemoryChannel; +import org.apache.flume.conf.Configurables; +import org.apache.flume.event.EventBuilder; +import org.bson.Document; +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.Test; + +/** + * Integration tests that exercise {@link MongoDbSink} against a real (but + * embedded/in-memory) MongoDB instance provided by Flapdoodle, verifying + * that duplicate key errors are tolerated instead of failing the whole + * batch and are tracked via a dedicated counter. + */ +public class TestMongoDbSinkEmbedded { + + private static MongodExecutable mongodExecutable; + private static MongoClient mongoClient; + private static int port; + + @BeforeClass + public static void startMongo() throws Exception { + port = Network.getFreeServerPort(); + MongodStarter starter = MongodStarter.getDefaultInstance(); + MongodConfig mongodConfig = MongodConfig.builder() + .version(Version.Main.V4_4) + .net(new Net(port, Network.localhostIsIPv6())) + .build(); + MongodExecutable executable = starter.prepare(mongodConfig); + MongodProcess process = executable.start(); + mongodExecutable = executable; + mongoClient = MongoClients.create("mongodb://localhost:" + port); + // Keep a reference so the process isn't garbage collected/stopped early. + assert process != null; + } + + @AfterClass + public static void stopMongo() { + if (mongoClient != null) { + mongoClient.close(); + } + if (mongodExecutable != null) { + mongodExecutable.stop(); + } + } + + private static Context baseContext(String database, String collection) { + Context context = new Context(); + context.put(MongoDbSinkConstants.CONNECTION_URI, "mongodb://localhost:" + port); + context.put(MongoDbSinkConstants.DATABASE_NAME, database); + context.put(MongoDbSinkConstants.COLLECTION, collection); + return context; + } + + private static MongoDbSink createAndStartSink(Context context) { + MongoDbSink sink = new MongoDbSink(); + Channel channel = new MemoryChannel(); + Configurables.configure(channel, new Context()); + sink.setChannel(channel); + channel.start(); + Configurables.configure(sink, context); + sink.start(); + return sink; + } + + private static void putEvent(Channel channel, String json) { + Transaction tx = channel.getTransaction(); + tx.begin(); + channel.put(EventBuilder.withBody(json.getBytes(StandardCharsets.UTF_8), new HashMap<>())); + tx.commit(); + tx.close(); + } + + @Test + public void testDuplicateKeyDoesNotFailBatchAndIsCountedSeparately() throws EventDeliveryException { + String database = "testDb1"; + String collectionName = "events"; + Context context = baseContext(database, collectionName); + MongoDbSink sink = createAndStartSink(context); + try { + MongoCollection collection = + mongoClient.getDatabase(database).getCollection(collectionName); + collection.createIndex(Indexes.ascending("uid"), new IndexOptions().unique(true)); + + Channel channel = sink.getChannel(); + // Two distinct events plus one that duplicates the first's unique key. + putEvent(channel, "{\"uid\":1,\"value\":\"a\"}"); + putEvent(channel, "{\"uid\":2,\"value\":\"b\"}"); + putEvent(channel, "{\"uid\":1,\"value\":\"c\"}"); + + Sink.Status status = sink.process(); + + assertEquals(Sink.Status.READY, status); + assertEquals(2, collection.countDocuments()); + assertEquals(1, sink.getDuplicateEventCount()); + } finally { + sink.stop(); + } + } + + @Test + public void testNoDuplicatesLeavesDuplicateCountAtZero() throws EventDeliveryException { + String database = "testDb2"; + String collectionName = "events"; + Context context = baseContext(database, collectionName); + MongoDbSink sink = createAndStartSink(context); + try { + MongoCollection collection = + mongoClient.getDatabase(database).getCollection(collectionName); + collection.createIndex(Indexes.ascending("uid"), new IndexOptions().unique(true)); + + Channel channel = sink.getChannel(); + putEvent(channel, "{\"uid\":1,\"value\":\"a\"}"); + putEvent(channel, "{\"uid\":2,\"value\":\"b\"}"); + + Sink.Status status = sink.process(); + + assertEquals(Sink.Status.READY, status); + assertEquals(2, collection.countDocuments()); + assertEquals(0, sink.getDuplicateEventCount()); + } finally { + sink.stop(); + } + } +} From f036eb9b2c706496de87cd8a18f2a437bc0aa6ec Mon Sep 17 00:00:00 2001 From: "Piotr P. Karwasz" Date: Wed, 2 Sep 2026 19:44:50 +0200 Subject: [PATCH 2/6] Use Docker for MongoDB integration tests --- flume-mongodb-sink/pom.xml | 93 +++++++++++++++++-- ...DbSinkEmbedded.java => MongoDbSinkIT.java} | 35 ++----- 2 files changed, 92 insertions(+), 36 deletions(-) rename flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/{TestMongoDbSinkEmbedded.java => MongoDbSinkIT.java} (77%) diff --git a/flume-mongodb-sink/pom.xml b/flume-mongodb-sink/pom.xml index b221f7e..a9b01ce 100644 --- a/flume-mongodb-sink/pom.xml +++ b/flume-mongodb-sink/pom.xml @@ -31,6 +31,9 @@ 4 0 org.apache.flume.sink.mongodb + + + 0.49.0 @@ -76,13 +79,6 @@ test - - de.flapdoodle.embed - de.flapdoodle.embed.mongo - 3.5.4 - test - - @@ -140,4 +136,87 @@ + + + docker + + + + linux + + + env.CI + true + + + + + + io.fabric8 + docker-maven-plugin + ${docker-maven-plugin.version} + + all + + + mongo + mongo:latest + + + + localhost:mongo.port:27017 + + + + + 27017 + + + + + + + + + + + start-mongo + + start + + + + stop-mongo + + stop + + + + + + org.apache.maven.plugins + maven-failsafe-plugin + + + + integration-test + verify + + + + + ${mongo.port} + + + + + + + + + + diff --git a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java similarity index 77% rename from flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java rename to flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java index cc4ffbf..a7e4f51 100644 --- a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/TestMongoDbSinkEmbedded.java +++ b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java @@ -23,13 +23,6 @@ import com.mongodb.client.MongoCollection; import com.mongodb.client.model.IndexOptions; import com.mongodb.client.model.Indexes; -import de.flapdoodle.embed.mongo.MongodExecutable; -import de.flapdoodle.embed.mongo.MongodProcess; -import de.flapdoodle.embed.mongo.MongodStarter; -import de.flapdoodle.embed.mongo.config.MongodConfig; -import de.flapdoodle.embed.mongo.config.Net; -import de.flapdoodle.embed.mongo.distribution.Version; -import de.flapdoodle.embed.process.runtime.Network; import java.nio.charset.StandardCharsets; import java.util.HashMap; import org.apache.flume.Channel; @@ -46,41 +39,25 @@ import org.junit.Test; /** - * Integration tests that exercise {@link MongoDbSink} against a real (but - * embedded/in-memory) MongoDB instance provided by Flapdoodle, verifying - * that duplicate key errors are tolerated instead of failing the whole - * batch and are tracked via a dedicated counter. + * Integration tests that exercise {@link MongoDbSink} against MongoDB running + * in the Docker container configured by the Maven {@code docker} profile. */ -public class TestMongoDbSinkEmbedded { +public class MongoDbSinkIT { - private static MongodExecutable mongodExecutable; private static MongoClient mongoClient; private static int port; @BeforeClass - public static void startMongo() throws Exception { - port = Network.getFreeServerPort(); - MongodStarter starter = MongodStarter.getDefaultInstance(); - MongodConfig mongodConfig = MongodConfig.builder() - .version(Version.Main.V4_4) - .net(new Net(port, Network.localhostIsIPv6())) - .build(); - MongodExecutable executable = starter.prepare(mongodConfig); - MongodProcess process = executable.start(); - mongodExecutable = executable; + public static void connectToMongo() { + port = Integer.parseInt(System.getProperty("mongo.port")); mongoClient = MongoClients.create("mongodb://localhost:" + port); - // Keep a reference so the process isn't garbage collected/stopped early. - assert process != null; } @AfterClass - public static void stopMongo() { + public static void closeMongoClient() { if (mongoClient != null) { mongoClient.close(); } - if (mongodExecutable != null) { - mongodExecutable.stop(); - } } private static Context baseContext(String database, String collection) { From 8829caa74cca063c4cb118509e1d0cbe12fc9a0f Mon Sep 17 00:00:00 2001 From: "Piotr P. Karwasz" Date: Wed, 2 Sep 2026 20:41:57 +0200 Subject: [PATCH 3/6] fix: remove unused `closed` field --- .../java/org/apache/flume/sink/mongodb/TestMongoDbSink.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) 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 a9a65c2..9d3ed4d 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 @@ -48,7 +48,6 @@ public class TestMongoDbSink { */ private static final class FakeMongoDbWriter implements MongoDbWriter { private final Map> written = new LinkedHashMap<>(); - private boolean closed = false; private long duplicateCountToReport = 0; @Override @@ -60,9 +59,7 @@ public MongoDbWriteResult write(String collectionName, List documents) } @Override - public void close() { - closed = true; - } + public void close() {} } private static Context baseContext() { From e38af70bee74ed335e83ad5cc4cab09e3e227d43 Mon Sep 17 00:00:00 2001 From: "Piotr P. Karwasz" Date: Wed, 2 Sep 2026 20:42:55 +0200 Subject: [PATCH 4/6] fix: remove `.flattened-pom.xml` --- .flattened-pom.xml | 216 --------------------------------------------- 1 file changed, 216 deletions(-) delete mode 100644 .flattened-pom.xml diff --git a/.flattened-pom.xml b/.flattened-pom.xml deleted file mode 100644 index 815ad8e..0000000 --- a/.flattened-pom.xml +++ /dev/null @@ -1,216 +0,0 @@ - - - - 4.0.0 - - org.apache.flume - flume-parent - 2.0.0-SNAPSHOT - - - org.apache.flume - flume-mongodb-parent - 2.0.0-SNAPSHOT - pom - Flume MongoDB Parent - https://logging.apache.org/flume/2.x/index.html/flume-mongodb-parent - 2022 - - Apache Software Foundation - http://www.apache.org - - - - The Apache Software License, Version 2.0 - http://www.apache.org/licenses/LICENSE-2.0.txt - - - - - rgoers - Ralph Goers - rgoers@apache.org - Intuit - - - - - Flume User List - user-subscribe@flume.apache.org - user-unsubscribe@flume.apache.org - user@flume.apache.org - http://mail-archives.apache.org/mod_mbox/flume-user/ - - - Flume Developer List - dev-subscribe@flume.apache.org - dev-unsubscribe@flume.apache.org - dev@flume.apache.org - http://mail-archives.apache.org/mod_mbox/flume-dev/ - - - Flume Commits - commits-subscribe@flume.apache.org - commits-unsubscribe@flume.apache.org - commits@flume.apache.org - http://mail-archives.apache.org/mod_mbox/flume-commits/ - - - - flume-mongodb-sink - - - https://gitbox.apache.org/repos/asf/flume-spring-boot.git - https://gitbox.apache.org/repos/asf/flume-spring-boot.git - https://gitbox.apache.org/repos/asf/flume-spring-boot.git - - - JIRA - https://issues.apache.org/jira/browse/FLUME - - - 5.10.0 - 2.17.0 - 1.12.0 - 2.0.0 - Ralph Goers - 0.12 - 2.26.1 - 4.13.2 - B3D8E1BA - 11 - 2.27.2 - 1.6 - 1.11 - UTF-8 - org.apache.flume.mongodb - 2.9 - 1.9.0 - 4.1.18 - 2.0.0-SNAPSHOT - 4.7.2.1 - 11 - rgoers@apache.org - - - - - org.apache.flume - flume-ng-core - ${flume.version} - - - org.apache.flume - flume-ng-sdk - ${flume.version} - - - org.apache.flume - flume-ng-configuration - ${flume.version} - - - org.mongodb - mongodb-driver-bom - ${mongodb.version} - pom - import - - - io.dropwizard.metrics - metrics-core - ${dropwizard-metrics.version} - - - org.apache.logging.log4j - log4j-api - ${log4j.version} - - - org.apache.logging.log4j - log4j-core - ${log4j.version} - - - com.fasterxml.jackson.core - jackson-core - ${jackson.version} - - - junit - junit - ${junit.version} - - - org.mockito - mockito-all - ${mockito.version} - test - - - - - - - org.apache.rat - apache-rat-plugin - ${rat.version} - - - verify.rat - verify - - check - - - - - - **/.idea/ - **/*.iml - src/main/resources/META-INF/services/**/* - **/nb-configuration.xml - .git/ - patchprocess/ - .gitignore - **/*.yml - **/*.yaml - **/*.json - .repository/ - **/*.diff - **/*.patch - **/*.avsc - **/*.avro - **/docs/** - **/test/resources/** - **/.settings/* - **/.classpath - **/.project - **/target/** - **/derby.log - **/metastore_db/ - .mvn/** - **/exclude-pmd.properties - - true - - - - - From 49a810404204342fb41607ba2a7dd342b5530e44 Mon Sep 17 00:00:00 2001 From: Ralph Goers Date: Sat, 5 Sep 2026 16:07:29 -0700 Subject: [PATCH 5/6] Use TestContainers --- flume-mongodb-sink/pom.xml | 119 +++++------------- .../flume/sink/mongodb/MongoDbSinkIT.java | 29 +++-- pom.xml | 9 ++ 3 files changed, 65 insertions(+), 92 deletions(-) diff --git a/flume-mongodb-sink/pom.xml b/flume-mongodb-sink/pom.xml index 48c2117..5d0106f 100644 --- a/flume-mongodb-sink/pom.xml +++ b/flume-mongodb-sink/pom.xml @@ -33,9 +33,6 @@ 4 0 org.apache.flume.sink.mongodb - - - 0.49.0 @@ -70,6 +67,12 @@ test + + org.testcontainers + mongodb + test + + org.apache.logging.log4j log4j-api @@ -81,88 +84,34 @@ test + + org.apache.logging.log4j + log4j-slf4j2-impl + test + + - - - docker - - - - linux - - - env.CI - true - - - - - - io.fabric8 - docker-maven-plugin - ${docker-maven-plugin.version} - - all - - - mongo - mongo:latest - - - - localhost:mongo.port:27017 - - - - - 27017 - - - - - - - - - - - start-mongo - - start - - - - stop-mongo - - stop - - - - - - org.apache.maven.plugins - maven-failsafe-plugin - - - - integration-test - verify - - - - - ${mongo.port} - - - - - - - - - + + + + + + org.apache.maven.plugins + maven-failsafe-plugin + + + + integration-test + verify + + + + + + diff --git a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java index a7e4f51..438f0cf 100644 --- a/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java +++ b/flume-mongodb-sink/src/test/java/org/apache/flume/sink/mongodb/MongoDbSinkIT.java @@ -35,34 +35,49 @@ import org.apache.flume.event.EventBuilder; import org.bson.Document; import org.junit.AfterClass; +import org.junit.Assume; import org.junit.BeforeClass; import org.junit.Test; +import org.testcontainers.DockerClientFactory; +import org.testcontainers.containers.MongoDBContainer; +import org.testcontainers.utility.DockerImageName; /** * Integration tests that exercise {@link MongoDbSink} against MongoDB running - * in the Docker container configured by the Maven {@code docker} profile. + * in a Testcontainers-managed Docker container. The tests are skipped + * (rather than failed) when Docker is not available in the environment. */ public class MongoDbSinkIT { + private static final DockerImageName MONGO_IMAGE = DockerImageName.parse("mongo:latest"); + + private static MongoDBContainer mongoDBContainer; private static MongoClient mongoClient; - private static int port; @BeforeClass - public static void connectToMongo() { - port = Integer.parseInt(System.getProperty("mongo.port")); - mongoClient = MongoClients.create("mongodb://localhost:" + port); + public static void startMongoContainer() { + Assume.assumeTrue( + "Docker is not available, skipping MongoDbSinkIT", + DockerClientFactory.instance().isDockerAvailable()); + + mongoDBContainer = new MongoDBContainer(MONGO_IMAGE); + mongoDBContainer.start(); + mongoClient = MongoClients.create(mongoDBContainer.getConnectionString()); } @AfterClass - public static void closeMongoClient() { + public static void stopMongoContainer() { if (mongoClient != null) { mongoClient.close(); } + if (mongoDBContainer != null) { + mongoDBContainer.stop(); + } } private static Context baseContext(String database, String collection) { Context context = new Context(); - context.put(MongoDbSinkConstants.CONNECTION_URI, "mongodb://localhost:" + port); + context.put(MongoDbSinkConstants.CONNECTION_URI, mongoDBContainer.getConnectionString()); context.put(MongoDbSinkConstants.DATABASE_NAME, database); context.put(MongoDbSinkConstants.COLLECTION, collection); return context; diff --git a/pom.xml b/pom.xml index 9a2fb4b..65e10d5 100644 --- a/pom.xml +++ b/pom.xml @@ -55,6 +55,7 @@ 4.1.18 5.10.0 + 1.21.4 @@ -74,6 +75,14 @@ pom import + + + org.testcontainers + testcontainers-bom + ${testcontainers.version} + pom + import + From e5c54c6366ba31e4082b4a939df831a5b4d1e43b Mon Sep 17 00:00:00 2001 From: Ralph Goers Date: Sat, 5 Sep 2026 17:08:01 -0700 Subject: [PATCH 6/6] Attempt to write the whole batch and only revert to one and a time when an error occurs --- .../sink/mongodb/DefaultMongoDbWriter.java | 56 ++++++++++++++++++- 1 file changed, 54 insertions(+), 2 deletions(-) diff --git a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java index 28fb4ab..6e63c6f 100644 --- a/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java +++ b/flume-mongodb-sink/src/main/java/org/apache/flume/sink/mongodb/DefaultMongoDbWriter.java @@ -18,8 +18,10 @@ import com.mongodb.DuplicateKeyException; import com.mongodb.ErrorCategory; +import com.mongodb.MongoBulkWriteException; import com.mongodb.MongoWriteException; import com.mongodb.WriteConcern; +import com.mongodb.bulk.BulkWriteError; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; import java.util.List; @@ -50,10 +52,60 @@ public MongoDbWriteResult write(String collectionName, List documents) collection = collection.withWriteConcern(writeConcern); } + // First try to insert the whole batch in a single operation, which is + // far more efficient than one insert per document. Only fall back to + // inserting one document at a time if the batch insert fails because + // of duplicate keys. + try { + collection.insertMany(documents); + return new MongoDbWriteResult(documents.size(), 0); + } catch (MongoBulkWriteException ex) { + if (isDuplicateKeyOnly(ex)) { + // insertMany() is ordered by default, so it stops at the first + // failing document; everything before that point was already + // successfully persisted. Skip those already-written documents + // before retrying the remainder one at a time, otherwise they + // would be re-attempted and incorrectly counted as duplicates. + int alreadyInserted = ex.getWriteResult().getInsertedCount(); + logger.warn( + "Duplicate key(s) while batch inserting into collection {}, " + + "retrying remaining documents one at a time: {}", + collectionName, + ex.getMessage()); + MongoDbWriteResult retryResult = writeOneAtATime( + collection, collectionName, documents.subList(alreadyInserted, documents.size())); + return new MongoDbWriteResult( + alreadyInserted + retryResult.getInsertedCount(), retryResult.getDuplicateCount()); + } + throw ex; + } + } + + /** + * Returns {@code true} if every error reported by the bulk write failure + * is a duplicate key error. + */ + private boolean isDuplicateKeyOnly(MongoBulkWriteException ex) { + List errors = ex.getWriteErrors(); + if (errors.isEmpty()) { + return false; + } + for (BulkWriteError error : errors) { + if (ErrorCategory.fromErrorCode(error.getCode()) != ErrorCategory.DUPLICATE_KEY) { + return false; + } + } + return true; + } + + /** + * Inserts documents one at a time so that a duplicate key on any single + * document does not prevent the rest of the batch from being inserted. + */ + private MongoDbWriteResult writeOneAtATime( + MongoCollection collection, String collectionName, List documents) { long insertedCount = 0; long duplicateCount = 0; - // Insert one document at a time (rather than insertMany) so that a - // single duplicate key does not abort the rest of the batch. for (Document document : documents) { try { collection.insertOne(document);