From fb9c350c3346f7804fffad91176af8d0cc9269d5 Mon Sep 17 00:00:00 2001
From: Ralph Goers
Date: Sat, 20 Jun 2026 06:39:08 -0700
Subject: [PATCH 1/7] Add RoutableProxyChannelSelector and allow Interceptors
to be configured in Spring Boot
---
.../conf/channel/ChannelSelectorType.java | 6 +-
flume-ng-core/pom.xml | 5 +
.../flume/channel/ChannelProcessor.java | 7 ++
.../channel/LoadBalancingChannelSelector.java | 2 +
.../channel/MultiplexingChannelSelector.java | 2 +
.../channel/RoutableProxyChannelSelector.java | 119 ++++++++++++++++++
.../flume/lifecycle/LifecycleSupervisor.java | 2 +
.../sink/AbstractSingleSinkProcessor.java | 2 +
.../flume/sink/AbstractSinkProcessor.java | 2 +
.../flume/sink/FailoverSinkProcessor.java | 2 +
.../sink/LoadBalancingSinkProcessor.java | 2 +
.../org/apache/flume/source/ExecSource.java | 2 +
.../TestRoutableProxyChannelSelector.java | 111 ++++++++++++++++
flume-parent/pom.xml | 6 +-
14 files changed, 266 insertions(+), 4 deletions(-)
create mode 100644 flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
create mode 100644 flume-ng-core/src/test/java/org/apache/flume/channel/TestRoutableProxyChannelSelector.java
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelSelectorType.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelSelectorType.java
index 37bc1bcde5..c15a8a79b7 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelSelectorType.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelSelectorType.java
@@ -41,7 +41,11 @@ public enum ChannelSelectorType implements ComponentWithClassName {
/**
* Multiplexing channel selector.
*/
- MULTIPLEXING("org.apache.flume.channel.MultiplexingChannelSelector");
+ MULTIPLEXING("org.apache.flume.channel.MultiplexingChannelSelector"),
+ /**
+ * Routable proxy channel selector.
+ */
+ ROUTABLE_PROXY("org.apache.flume.channel.RoutableProxyChannelSelector");
private final String channelSelectorClassName;
diff --git a/flume-ng-core/pom.xml b/flume-ng-core/pom.xml
index a41891d265..c6429bba24 100644
--- a/flume-ng-core/pom.xml
+++ b/flume-ng-core/pom.xml
@@ -144,6 +144,11 @@
mockito-core
test
+
+ com.github.spotbugs
+ spotbugs-annotations
+ provided
+
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/ChannelProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/channel/ChannelProcessor.java
index c3d85ed4a2..30ef43b544 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/ChannelProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/ChannelProcessor.java
@@ -55,8 +55,15 @@ public class ChannelProcessor implements Configurable {
private final InterceptorChain interceptorChain;
public ChannelProcessor(ChannelSelector selector) {
+ this(selector, null);
+ }
+
+ public ChannelProcessor(ChannelSelector selector, List interceptors) {
this.selector = selector;
this.interceptorChain = new InterceptorChain();
+ if (interceptors != null) {
+ interceptorChain.setInterceptors(interceptors);
+ }
}
public void initialize() {
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
index 5aea76d28b..cb2b9242ce 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
@@ -18,6 +18,7 @@
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Collections;
import java.util.List;
import java.util.Random;
@@ -39,6 +40,7 @@
* defaults to ROUND_ROBIN type, but can be overridden via
* configuration.
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class LoadBalancingChannelSelector extends AbstractChannelSelector {
private final List emptyList = Collections.emptyList();
private ChannelPicker picker;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
index 2a500980e2..64ea9276b6 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
@@ -16,6 +16,7 @@
*/
package org.apache.flume.channel;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@@ -25,6 +26,7 @@
import org.apache.flume.Event;
import org.apache.flume.FlumeException;
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class MultiplexingChannelSelector extends AbstractChannelSelector {
public static final String CONFIG_MULTIPLEX_HEADER_NAME = "header";
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
new file mode 100644
index 0000000000..8fcda6c47c
--- /dev/null
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
@@ -0,0 +1,119 @@
+/*
+ * 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.channel;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.flume.Channel;
+import org.apache.flume.ChannelSelector;
+import org.apache.flume.Context;
+import org.apache.flume.Event;
+import org.apache.flume.conf.Configurables;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+public class RoutableProxyChannelSelector extends LoadBalancingChannelSelector {
+ private static final Logger LOGGER = LogManager.getLogger(RoutableProxyChannelSelector.class);
+ private static final String HEADER_NAME = "headerName";
+ private static final String SELECTOR = ".selector.";
+ private static final String CHANNELS = "channels";
+ private static final String DEFAULT = "default";
+ private static final String TYPE = "type";
+
+ private final Map selectorMap = new HashMap<>();
+ private ChannelSelector defaultSelector;
+ private String headerName;
+
+ public void addSelector(String headerName, ChannelSelector selector) {
+ selectorMap.put(headerName, selector);
+ }
+
+ public void setDefaultSelector(ChannelSelector defaultSelector) {
+ this.defaultSelector = defaultSelector;
+ }
+
+ public ChannelSelector getDefaultSelector() {
+ return defaultSelector;
+ }
+
+ @Override
+ public void configure(Context context) {
+ Configurables.ensureRequiredNonNull(context, HEADER_NAME);
+ List allChannels = getAllChannels();
+ for (Map.Entry entry : context.getParameters().entrySet()) {
+ if (entry.getKey().equals(HEADER_NAME)) {
+ this.headerName = entry.getValue();
+ } else if (!entry.getKey().equals(TYPE)) {
+ String key = StringUtils.substringBefore(entry.getKey(), ".");
+ Map map = context.getSubProperties(key + SELECTOR);
+ if (map != null) {
+ String channelNames = getRequiredNonNull(map, CHANNELS, key);
+ Set channelSet = new HashSet<>(Arrays.asList(channelNames.split("\\s+")));
+ List channels = new ArrayList<>(channelSet.size());
+ for (String channelName : channelSet) {
+ for (Channel channel : allChannels) {
+ if (channelName.equals(channel.getName())) {
+ channels.add(channel);
+ }
+ }
+ }
+ ChannelSelector selector = ChannelSelectorFactory.create(channels, map);
+ if (DEFAULT.equals(key)) {
+ defaultSelector = selector;
+ } else {
+ selectorMap.put(key, selector);
+ }
+ }
+ }
+ }
+ if (headerName == null) {
+ throw new IllegalArgumentException("No header name specified for RoutableProxy");
+ }
+ }
+
+ @Override
+ public List getRequiredChannels(Event event) {
+ return getSelector(event).getRequiredChannels(event);
+ }
+
+ @Override
+ public List getOptionalChannels(Event event) {
+ return getSelector(event).getOptionalChannels(event);
+ }
+
+ private ChannelSelector getSelector(Event event) {
+ ChannelSelector channelSelector = selectorMap.get(event.getHeaders().get(headerName));
+ if (channelSelector == null) {
+ channelSelector = defaultSelector;
+ }
+ return channelSelector;
+ }
+
+ private String getRequiredNonNull(Map map, String keyName, String prefix) {
+ String value = map.get(keyName);
+ if (value == null) {
+ throw new IllegalArgumentException(String.format("Missing key %s in %s", keyName, prefix));
+ }
+ return value;
+ }
+}
diff --git a/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java b/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
index 67ed39acfb..8151c0c6fd 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
@@ -18,6 +18,7 @@
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.HashMap;
import java.util.Map;
import java.util.Map.Entry;
@@ -29,6 +30,7 @@
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class LifecycleSupervisor implements LifecycleAware {
private static final Logger logger = LogManager.getLogger();
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
index 343ee73b35..d3f453b3ec 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
@@ -17,6 +17,7 @@
package org.apache.flume.sink;
import com.google.common.base.Preconditions;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.List;
import org.apache.flume.Sink;
import org.apache.flume.SinkProcessor;
@@ -25,6 +26,7 @@
/**
* A Sink Processor that only accesses a single Sink.
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public abstract class AbstractSingleSinkProcessor implements SinkProcessor {
protected Sink sink;
private LifecycleState lifecycleState;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
index 2be6b16062..e57cfa008d 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
@@ -16,6 +16,7 @@
*/
package org.apache.flume.sink;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
@@ -26,6 +27,7 @@
/**
* A convenience base class for sink processors.
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public abstract class AbstractSinkProcessor implements SinkProcessor {
private LifecycleState state;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
index a09a3107ac..49961a30dc 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
@@ -16,6 +16,7 @@
*/
package org.apache.flume.sink;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -59,6 +60,7 @@
* host1.sinkgroups.group1.processor.maxpenalty = 10000
*
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class FailoverSinkProcessor extends AbstractSinkProcessor {
private static final int FAILURE_PENALTY = 1000;
private static final int DEFAULT_MAX_PENALTY = 30000;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
index 30943ca4e9..ad6f0609f8 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
@@ -17,6 +17,7 @@
package org.apache.flume.sink;
import com.google.common.base.Preconditions;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Iterator;
import java.util.List;
import org.apache.flume.Context;
@@ -75,6 +76,7 @@
* @see FailoverSinkProcessor
* @see LoadBalancingSinkProcessor.SinkSelector
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class LoadBalancingSinkProcessor extends AbstractSinkProcessor {
public static final String CONFIG_SELECTOR = "selector";
public static final String CONFIG_SELECTOR_PREFIX = CONFIG_SELECTOR + ".";
diff --git a/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java b/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
index 3a6778a2f2..2b11df2f06 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
@@ -18,6 +18,7 @@
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
+import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
@@ -141,6 +142,7 @@
* TODO
*
*/
+@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class ExecSource extends AbstractSource implements EventDrivenSource, Configurable, BatchSizeSupported {
private static final Logger logger = LogManager.getLogger();
diff --git a/flume-ng-core/src/test/java/org/apache/flume/channel/TestRoutableProxyChannelSelector.java b/flume-ng-core/src/test/java/org/apache/flume/channel/TestRoutableProxyChannelSelector.java
new file mode 100644
index 0000000000..f1df71cbf3
--- /dev/null
+++ b/flume-ng-core/src/test/java/org/apache/flume/channel/TestRoutableProxyChannelSelector.java
@@ -0,0 +1,111 @@
+/*
+ * 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.channel;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+import junit.framework.Assert;
+import org.apache.flume.Channel;
+import org.apache.flume.ChannelSelector;
+import org.apache.flume.Event;
+import org.apache.flume.event.SimpleEvent;
+import org.junit.Test;
+
+public class TestRoutableProxyChannelSelector {
+
+ private List channels = new ArrayList();
+ private static final String[] config = new String[] {
+ "headerName = processingMode",
+ "type = routable_proxy",
+ "normal.selector.type = load_balancing",
+ "normal.selector.channels = ch1 ch2",
+ "normal.selector.policy = round_robin",
+ "default.selector.type = load_balancing",
+ "default.selector.channels = ch3 ch4",
+ "default.selector.policy = round_robin"
+ };
+
+ private ChannelSelector selector;
+
+ @Test
+ public void testProxySelector() throws Exception {
+ channels.clear();
+ channels.add(MockChannel.createMockChannel("ch1"));
+ channels.add(MockChannel.createMockChannel("ch2"));
+ channels.add(MockChannel.createMockChannel("ch3"));
+ channels.add(MockChannel.createMockChannel("ch4"));
+ RoutableProxyChannelSelector selector =
+ (RoutableProxyChannelSelector) ChannelSelectorFactory.create(channels, getConfig());
+ Assert.assertNotNull(selector);
+ Assert.assertNotNull(selector.getDefaultSelector());
+ Event event = new SimpleEvent();
+ event.getHeaders().put("processingMode", "normal");
+ List channels = selector.getRequiredChannels(event);
+ Assert.assertNotNull(channels);
+ Assert.assertEquals(1, channels.size());
+ String channelName = channels.get(0).getName();
+ Assert.assertTrue(channelName.equals("ch1") || channelName.equals("ch2"));
+ }
+
+ @Test
+ public void testProxySelectorManualConfig() throws Exception {
+ channels.clear();
+ channels.add(MockChannel.createMockChannel("ch1"));
+ channels.add(MockChannel.createMockChannel("ch2"));
+ channels.add(MockChannel.createMockChannel("ch3"));
+ channels.add(MockChannel.createMockChannel("ch4"));
+ Map config = new HashMap<>();
+ config.put("headerName", "processingMode");
+ config.put("type", "routable_proxy");
+ RoutableProxyChannelSelector selector =
+ (RoutableProxyChannelSelector) ChannelSelectorFactory.create(channels, config);
+ config.clear();
+ config.put("policy", "round_robin");
+ config.put("type", "load_balancing");
+ List channels1 = new ArrayList<>();
+ channels1.add(channels.get(0));
+ channels1.add(channels.get(1));
+ LoadBalancingChannelSelector loadBalancingSelector =
+ (LoadBalancingChannelSelector) ChannelSelectorFactory.create(channels1, config);
+ selector.setDefaultSelector(loadBalancingSelector);
+ List channels2 = new ArrayList<>();
+ channels2.add(channels.get(2));
+ channels2.add(channels.get(3));
+ loadBalancingSelector = (LoadBalancingChannelSelector) ChannelSelectorFactory.create(channels2, config);
+ selector.addSelector("validation", loadBalancingSelector);
+ Assert.assertNotNull(selector);
+ Assert.assertNotNull(selector.getDefaultSelector());
+ Event event = new SimpleEvent();
+ event.getHeaders().put("processingMode", "normal");
+ List channels = selector.getRequiredChannels(event);
+ Assert.assertNotNull(channels);
+ Assert.assertEquals(1, channels.size());
+ String channelName = channels.get(0).getName();
+ Assert.assertTrue(channelName.equals("ch1") || channelName.equals("ch2"));
+ }
+
+ private Map getConfig() {
+ return Arrays.stream(config)
+ .map(line -> line.split("=", 2)) // Limit to 2 parts to keep values with '=' intact
+ .filter(parts -> parts.length == 2)
+ .collect(Collectors.toMap(parts -> parts[0].trim(), parts -> parts[1].trim()));
+ }
+}
diff --git a/flume-parent/pom.xml b/flume-parent/pom.xml
index 261737b658..57bc9d9c09 100644
--- a/flume-parent/pom.xml
+++ b/flume-parent/pom.xml
@@ -237,9 +237,9 @@
1.15.0
5.9.0
10.17.1.0
- 2.19.0
- 2.21.1
- 2.21.1
+ 2.22
+ 2.22.0
+ 2.22.0
1.4.1
2.14.0
33.4.8-jre
From 5ca1c953bbb30e061f64faf76204364269d4f1cb Mon Sep 17 00:00:00 2001
From: Ralph Goers
Date: Thu, 25 Jun 2026 07:56:57 -0700
Subject: [PATCH 2/7] Enhancements for Spring Boot
---
flume-bom/pom.xml | 8 +-
flume-ng-channels/flume-file-channel/pom.xml | 1 -
.../apache/flume/conf/FlumeConfiguration.java | 2 +
.../flume/conf/channel/ChannelType.java | 14 +
flume-ng-core/pom.xml | 1 -
flume-ng-node/pom.xml | 17 -
flume-parent/pom.xml | 427 +-------------
flume-third-party/pom.xml | 527 ++++++++++++++++++
pom.xml | 1 +
9 files changed, 559 insertions(+), 439 deletions(-)
create mode 100644 flume-third-party/pom.xml
diff --git a/flume-bom/pom.xml b/flume-bom/pom.xml
index 67ce79ed2e..67f840d200 100644
--- a/flume-bom/pom.xml
+++ b/flume-bom/pom.xml
@@ -38,11 +38,11 @@
2.0.0-SNAPSHOT
2.0.0-SNAPSHOT
2.0.0-SNAPSHOT
- 2.0.0-SNAPSHOT
+ 2.1.0-SNAPSHOT
2.0.0-SNAPSHOT
2.0.0-SNAPSHOT
2.0.0-SNAPSHOT
- 2.0.0-SNAPSHOT
+ 2.0.0-SNAPSHOT
@@ -105,12 +105,12 @@
org.apache.flume
flume-rpc-avro
- ${project.version}
+ ${flume-rpc.version}
org.apache.flume
flume-rpc-thrift
- ${project.version}
+ ${flume-rpc.version}
org.apache.flume
diff --git a/flume-ng-channels/flume-file-channel/pom.xml b/flume-ng-channels/flume-file-channel/pom.xml
index d4052e921e..9e7e219c30 100644
--- a/flume-ng-channels/flume-file-channel/pom.xml
+++ b/flume-ng-channels/flume-file-channel/pom.xml
@@ -121,7 +121,6 @@
com.google.code.findbugs
jsr305
- ${jsr305.version}
provided
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
index 1fd692a0ce..a7bc35fe39 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
@@ -69,6 +69,7 @@
import org.apache.flume.configfilter.ConfigFilter;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
+import org.apache.logging.log4j.internal.annotation.SuppressFBWarnings;
/**
*
@@ -83,6 +84,7 @@
* @see org.apache.flume.node.ConfigurationProvider
*
*/
+@SuppressFBWarnings(value = {"EI_EXPOSE_REP"})
public class FlumeConfiguration {
private static final Logger logger = LogManager.getLogger();
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
index e782472ce1..315c24c5ac 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
@@ -46,6 +46,11 @@ public enum ChannelType implements ComponentWithClassName {
*/
JDBC("org.apache.flume.channel.jdbc.JdbcChannel"),
+ /**
+ * Kafka channel.
+ */
+ KAFKA("org.apache.flume.channel.kafka.KafkaChannel"),
+
/**
* Spillable Memory channel
*
@@ -59,6 +64,15 @@ private ChannelType(String channelClassName) {
this.channelClassName = channelClassName;
}
+ private static ChannelType getChannelTypeByName(String type) {
+ for (ChannelType channelType : ChannelType.values()) {
+ if (channelType.channelClassName.equals(type)) {
+ return channelType;
+ }
+ }
+ return null;
+ }
+
@Deprecated
public String getChannelClassName() {
return channelClassName;
diff --git a/flume-ng-core/pom.xml b/flume-ng-core/pom.xml
index c6429bba24..7bb547389d 100644
--- a/flume-ng-core/pom.xml
+++ b/flume-ng-core/pom.xml
@@ -89,7 +89,6 @@
com.google.code.findbugs
jsr305
- ${jsr305.version}
provided
diff --git a/flume-ng-node/pom.xml b/flume-ng-node/pom.xml
index c746cc4f68..8ac7986bb3 100644
--- a/flume-ng-node/pom.xml
+++ b/flume-ng-node/pom.xml
@@ -54,22 +54,6 @@
flume-ng-core
-
-
org.apache.flume.flume-ng-channels
flume-file-channel
@@ -120,7 +104,6 @@
com.google.code.findbugs
jsr305
- ${jsr305.version}
provided
diff --git a/flume-parent/pom.xml b/flume-parent/pom.xml
index 57bc9d9c09..ddebdd3638 100644
--- a/flume-parent/pom.xml
+++ b/flume-parent/pom.xml
@@ -209,6 +209,7 @@
+ 2.0.0-SNAPSHOT
2022-01-02T00:00:00Z
@@ -227,34 +228,10 @@
${project.basedir}/target/docs
2.3.7
- 1.5.0
- 1.22.0
- 3.2.2
- 1.28.0
- 2.22.0
- 3.20.0
- 1.4.0
- 1.15.0
5.9.0
- 10.17.1.0
- 2.22
- 2.22.0
- 2.22.0
1.4.1
- 2.14.0
- 33.4.8-jre
3.5.0
- 4.4.15
- 1.84
- 4.5.13
- 1.10
- 6.1.0
- 12.1.9
4.13.2
- 3.0.2
- 2.26.0
- 0.9.9
- 2.2.8
5.18.0
1.8
@@ -272,31 +249,23 @@
4.9.3
${spotbugs.version}.0
3.5.5
- 4.2.15.Final
4.35.0
0.6.1
- 1.7.0
0.12
- 2.0.17
- 1.1.10.8
- 4.9.3
- 1.19.0
1.53
- 1.1.3
- 3.8.6
1.7.1
true
3.4.0
+ 1.19.0
2.90.0
-
- org.bouncycastle
- bc-jdk18on-bom
- ${bouncycastle.version}
+ org.apache.flume
+ flume-third-party-dependencies
+ ${flume-project.version}
pom
import
@@ -328,12 +297,6 @@
test
-
- org.apache.flume
- flume-rpc-avro
- ${project.version}
-
-
org.apache.hadoop
hadoop-minikdc
@@ -364,185 +327,11 @@
test
-
-
-
- commons-cli
- commons-cli
- ${commons-cli.version}
-
-
-
-
- commons-logging
- commons-logging
- ${commons-logging.version}
-
-
-
- org.apache.commons
- commons-lang3
- ${commons-lang.version}
-
-
-
- org.apache.commons
- commons-text
- ${commons-text.version}
-
-
-
- com.google.guava
- guava
- ${guava.version}
-
-
-
- org.apache.logging.log4j
- log4j-bom
- ${log4j.version}
- pom
- import
-
-
-
-
- org.slf4j
- slf4j-bom
- ${slf4j.version}
- pom
- import
-
-
-
- com.google.protobuf
- protobuf-java
- ${external.protobuf.version}
- compile
-
-
-
- jakarta.servlet
- jakarta.servlet-api
-
- ${jakarta-servlet.version}
- provided
-
-
-
- org.eclipse.jetty.ee11
- jetty-ee11-servlet
- ${jetty.version}
-
-
-
- org.eclipse.jetty
- jetty-security
- ${jetty.version}
-
-
-
- org.eclipse.jetty
- jetty-util
- ${jetty.version}
-
-
-
- org.eclipse.jetty
- jetty-server
- ${jetty.version}
-
-
-
- org.eclipse.jetty
- jetty-jmx
- ${jetty.version}
-
-
-
- org.apache.httpcomponents
- httpclient
- ${httpclient.version}
-
-
-
- org.apache.httpcomponents
- httpcore
- ${httpcore.version}
-
-
-
- org.mapdb
- mapdb
- ${mapdb.version}
-
-
-
-
- com.google.code.gson
- gson
- ${gson.version}
-
-
-
- commons-codec
- commons-codec
- ${commons-codec.version}
-
-
-
- commons-io
- commons-io
- ${commons-io.version}
-
-
-
- commons-collections
- commons-collections
- ${commons-collections.version}
-
-
-
- org.apache.derby
- derby
- ${derby.version}
-
-
-
- com.fasterxml.jackson.core
- jackson-annotations
- ${jackson-annotations.version}
-
-
-
- com.fasterxml.jackson.core
- jackson-core
- ${fasterxml.jackson.version}
-
-
-
- com.fasterxml.jackson.core
- jackson-databind
- ${fasterxml.jackson.databind.version}
-
-
-
- org.schwering
- irclib
- ${irclib.version}
-
-
-
- com.jcraft
- jzlib
- ${zlib.version}
-
-
org.apache.flume
flume-dependencies
- ${project.version}
+ ${flume-project.version}
pom
import
@@ -550,66 +339,17 @@
org.apache.flume
build-support
- ${project.version}
+ ${flume-project.version}
org.apache.flume
flume-ng-sdk
- ${project.version}
+ ${flume-project.version}
tests
test
-
- org.apache.commons
- commons-compress
- ${commons-compress.version}
-
-
-
- org.apache.mina
- mina-core
- ${mina.version}
-
-
-
- io.netty
- netty-all
- ${netty-all.version}
-
-
-
- io.prometheus
- prometheus-metrics-core
- ${prometheus.version}
-
-
-
- io.prometheus
- prometheus-metrics-exporter-servlet-jakarta
- ${prometheus.version}
-
-
-
- org.xerial.snappy
- snappy-java
- ${snappy-java.version}
-
-
-
-
- org.apache.curator
- curator-framework
- ${curator.version}
-
-
-
- org.apache.curator
- curator-recipes
- ${curator.version}
-
-
org.apache.curator
curator-test
@@ -617,64 +357,6 @@
test
-
- org.apache.hadoop
- hadoop-common
- ${hadoop.version}
- true
-
-
- tomcat
- jasper-compiler
-
-
- tomcat
- jasper-runtime
-
-
- ch.qos.reload4j
- reload4j
-
-
- org.slf4j
- slf4j-reload4j
-
-
-
-
- org.apache.hadoop
- hadoop-hdfs
- ${hadoop.version}
-
-
- tomcat
- jasper-compiler
-
-
- tomcat
- jasper-runtime
-
-
- ch.qos.reload4j
- reload4j
-
-
-
-
- org.apache.hadoop
- hadoop-hdfs-client
- ${hadoop.version}
-
-
- tomcat
- jasper-compiler
-
-
- tomcat
- jasper-runtime
-
-
-
org.apache.hadoop
hadoop-minicluster
@@ -699,31 +381,6 @@
-
- org.apache.hadoop
- hadoop-client
- ${hadoop.version}
-
-
- org.apache.hadoop
- hadoop-annotations
- ${hadoop.version}
-
-
- org.apache.hadoop
- hadoop-auth
- ${hadoop.version}
-
-
- ch.qos.reload4j
- reload4j
-
-
- org.slf4j
- slf4j-reload4j
-
-
-
org.apache.hadoop
hadoop-mapreduce-client-core
@@ -740,37 +397,6 @@
-
- org.apache.hadoop
- hadoop-distcp
- ${hadoop.version}
-
-
-
- org.apache.zookeeper
- zookeeper
- ${zookeeper.version}
-
-
- ch.qos.logback
- *
-
-
-
-
-
-
- com.google.code.findbugs
- jsr305
- ${jsr305.version}
-
-
-
- com.github.spotbugs
- spotbugs-annotations
- ${spotbugs-annotations.version}
-
-
org.mock-server
mockserver-netty
@@ -831,10 +457,10 @@
${project.name}
- ${project.version}
+ ${flume-project.version}
${project.organization.name}
${project.name}
- ${project.version}
+ ${flume-project.version}
${project.organization.name}
org.apache
${maven.compiler.source}
@@ -1202,38 +828,6 @@
org.apache.maven.plugins
maven-pmd-plugin
-
-
- org.codehaus.mojo
- flatten-maven-plugin
- 1.7.3
-
-
- flatten-bom
-
- flatten
-
- process-resources
- false
-
- true
- bom
-
- resolve
- resolve
- flatten
-
-
-
-
-
- flatten.clean
-
- clean
-
-
-
-
@@ -1270,6 +864,7 @@
org.slf4j:slf4j-api
org.apache.logging.log4j:log4j-core:*:jar:test
+ org.apache.logging.log4j:log4j-jcl:*:jar:test
org.apache.logging.log4j:log4j-slf4j2-impl:*:jar:test
diff --git a/flume-third-party/pom.xml b/flume-third-party/pom.xml
new file mode 100644
index 0000000000..c9a82e6665
--- /dev/null
+++ b/flume-third-party/pom.xml
@@ -0,0 +1,527 @@
+
+
+
+
+ 4.0.0
+
+
+ org.apache.flume
+ flume-project
+ 2.0.0-SNAPSHOT
+
+
+ org.apache.flume
+ flume-third-party-dependencies
+ pom
+
+ Apache Flume Third Party Dependencies
+
+ 2009
+
+
+ Apache Software Foundation
+ http://www.apache.org
+
+
+
+
+ The Apache Software License, Version 2.0
+ http://www.apache.org/licenses/LICENSE-2.0.txt
+
+
+
+
+ scm:git:http://git-wip-us.apache.org/repos/asf/flume.git
+ scm:git:https://git-wip-us.apache.org/repos/asf/flume.git
+ https://git-wip-us.apache.org/repos/asf?p=flume.git;a=tree;h=refs/heads/trunk;hb=trunk
+
+
+
+ 2.0.0-SNAPSHOT
+ 1.5.0
+ 1.22.0
+ 3.2.2
+ 1.28.0
+ 2.22.0
+ 3.20.0
+ 1.4.0
+ 1.15.0
+ 5.9.0
+ 10.17.1.0
+ 2.22
+ 2.22.0
+ 2.22.0
+ 2.14.0
+ 33.4.8-jre
+ 3.5.0
+ 4.4.15
+ 1.84
+ 4.5.13
+ 1.10
+ 6.1.0
+ 12.1.9
+ 3.0.2
+ 2.26.0
+ 0.9.9
+ 2.2.8
+ 4.9.3
+ 4.2.15.Final
+ 4.35.0
+ 1.7.0
+ 2.0.17
+ 1.1.10.8
+ 4.9.3
+ 1.19.0
+ 1.1.3
+ 3.8.6
+ 1.7.1
+
+
+
+
+
+
+ org.bouncycastle
+ bc-jdk18on-bom
+ ${bouncycastle.version}
+ pom
+ import
+
+
+
+
+
+ commons-cli
+ commons-cli
+ ${commons-cli.version}
+
+
+
+
+ commons-logging
+ commons-logging
+ ${commons-logging.version}
+
+
+
+ org.apache.commons
+ commons-lang3
+ ${commons-lang.version}
+
+
+
+ org.apache.commons
+ commons-text
+ ${commons-text.version}
+
+
+
+ com.google.guava
+ guava
+ ${guava.version}
+
+
+
+ org.apache.logging.log4j
+ log4j-bom
+ ${log4j.version}
+ pom
+ import
+
+
+
+
+ org.slf4j
+ slf4j-bom
+ ${slf4j.version}
+ pom
+ import
+
+
+
+ com.google.protobuf
+ protobuf-java
+ ${external.protobuf.version}
+ compile
+
+
+
+ jakarta.servlet
+ jakarta.servlet-api
+
+ ${jakarta-servlet.version}
+ provided
+
+
+
+ org.eclipse.jetty.ee11
+ jetty-ee11-servlet
+ ${jetty.version}
+
+
+
+ org.eclipse.jetty
+ jetty-security
+ ${jetty.version}
+
+
+
+ org.eclipse.jetty
+ jetty-util
+ ${jetty.version}
+
+
+
+ org.eclipse.jetty
+ jetty-server
+ ${jetty.version}
+
+
+
+ org.eclipse.jetty
+ jetty-jmx
+ ${jetty.version}
+
+
+
+ org.apache.httpcomponents
+ httpclient
+ ${httpclient.version}
+
+
+
+ org.apache.httpcomponents
+ httpcore
+ ${httpcore.version}
+
+
+
+ org.mapdb
+ mapdb
+ ${mapdb.version}
+
+
+
+
+ com.google.code.gson
+ gson
+ ${gson.version}
+
+
+
+ commons-codec
+ commons-codec
+ ${commons-codec.version}
+
+
+
+ commons-io
+ commons-io
+ ${commons-io.version}
+
+
+
+ commons-collections
+ commons-collections
+ ${commons-collections.version}
+
+
+
+ org.apache.derby
+ derby
+ ${derby.version}
+
+
+
+ com.fasterxml.jackson.core
+ jackson-annotations
+ ${jackson-annotations.version}
+
+
+
+ com.fasterxml.jackson.core
+ jackson-core
+ ${fasterxml.jackson.version}
+
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ ${fasterxml.jackson.databind.version}
+
+
+
+ org.schwering
+ irclib
+ ${irclib.version}
+
+
+
+ com.jcraft
+ jzlib
+ ${zlib.version}
+
+
+
+
+ org.apache.flume
+ flume-dependencies
+ ${flume-project.version}
+ pom
+ import
+
+
+
+ org.apache.flume
+ build-support
+ ${flume-project.version}
+
+
+
+ org.apache.commons
+ commons-compress
+ ${commons-compress.version}
+
+
+
+ org.apache.mina
+ mina-core
+ ${mina.version}
+
+
+
+ io.netty
+ netty-all
+ ${netty-all.version}
+
+
+
+ io.prometheus
+ prometheus-metrics-core
+ ${prometheus.version}
+
+
+
+ io.prometheus
+ prometheus-metrics-exporter-servlet-jakarta
+ ${prometheus.version}
+
+
+
+ org.xerial.snappy
+ snappy-java
+ ${snappy-java.version}
+
+
+
+
+ org.apache.curator
+ curator-framework
+ ${curator.version}
+
+
+
+ org.apache.curator
+ curator-recipes
+ ${curator.version}
+
+
+
+ org.apache.curator
+ curator-test
+ ${curator.version}
+ test
+
+
+
+ org.apache.hadoop
+ hadoop-common
+ ${hadoop.version}
+ true
+
+
+ tomcat
+ jasper-compiler
+
+
+ tomcat
+ jasper-runtime
+
+
+ ch.qos.reload4j
+ reload4j
+
+
+ org.slf4j
+ slf4j-reload4j
+
+
+
+
+ org.apache.hadoop
+ hadoop-hdfs
+ ${hadoop.version}
+
+
+ tomcat
+ jasper-compiler
+
+
+ tomcat
+ jasper-runtime
+
+
+ ch.qos.reload4j
+ reload4j
+
+
+
+
+ org.apache.hadoop
+ hadoop-hdfs-client
+ ${hadoop.version}
+
+
+ tomcat
+ jasper-compiler
+
+
+ tomcat
+ jasper-runtime
+
+
+
+
+ org.apache.hadoop
+ hadoop-minicluster
+ ${hadoop.version}
+ test
+
+
+ tomcat
+ jasper-compiler
+
+
+ tomcat
+ jasper-runtime
+
+
+ ch.qos.reload4j
+ reload4j
+
+
+ org.slf4j
+ slf4j-reload4j
+
+
+
+
+ org.apache.hadoop
+ hadoop-client
+ ${hadoop.version}
+
+
+ org.apache.hadoop
+ hadoop-annotations
+ ${hadoop.version}
+
+
+ org.apache.hadoop
+ hadoop-auth
+ ${hadoop.version}
+
+
+ ch.qos.reload4j
+ reload4j
+
+
+ org.slf4j
+ slf4j-reload4j
+
+
+
+
+ org.apache.hadoop
+ hadoop-distcp
+ ${hadoop.version}
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ ${zookeeper.version}
+
+
+ ch.qos.logback
+ *
+
+
+
+
+
+
+ com.google.code.findbugs
+ jsr305
+ ${jsr305.version}
+
+
+
+ com.github.spotbugs
+ spotbugs-annotations
+ ${spotbugs-annotations.version}
+
+
+
+
+
+
+
+
+
+ org.codehaus.mojo
+ flatten-maven-plugin
+ 1.7.3
+
+
+ flatten-bom
+
+ flatten
+
+ process-resources
+ false
+
+ true
+ bom
+
+ resolve
+ flatten
+
+
+
+
+
+ flatten.clean
+
+ clean
+
+
+
+
+
+
+
diff --git a/pom.xml b/pom.xml
index c13d27277d..1c2e2ca743 100644
--- a/pom.xml
+++ b/pom.xml
@@ -33,6 +33,7 @@
flume-bom
+ flume-third-party
flume-parent
flume-ng-core
flume-ng-configuration
From d6de49d1499c515bea137c86a0bf288e8e40cb37 Mon Sep 17 00:00:00 2001
From: Ralph Goers
Date: Fri, 26 Jun 2026 07:00:54 -0700
Subject: [PATCH 3/7] Add relative path
---
flume-third-party/pom.xml | 2 +-
1 file changed, 1 insertion(+), 1 deletion(-)
diff --git a/flume-third-party/pom.xml b/flume-third-party/pom.xml
index c9a82e6665..73029cb1d9 100644
--- a/flume-third-party/pom.xml
+++ b/flume-third-party/pom.xml
@@ -23,7 +23,7 @@
org.apache.flume
flume-project
2.0.0-SNAPSHOT
-
+ ../pom.xml
org.apache.flume
flume-third-party-dependencies
From 2c16a6ca34b8e3acc42bb1da70dfbefe2762db60 Mon Sep 17 00:00:00 2001
From: "Piotr P. Karwasz"
Date: Fri, 26 Jun 2026 17:31:25 +0200
Subject: [PATCH 4/7] Drop LGPL spotbugs-annotations in favor of internal
annotation
SpotBugs honors any annotation whose simple name is SuppressFBWarnings,
so a tiny internal annotation lets us avoid shipping a build dependency
on the LGPL-licensed spotbugs-annotations artifact.
Assisted-By: Claude Opus 4.8 (1M context)
---
.../apache/flume/conf/FlumeConfiguration.java | 2 +-
.../conf/internal/SuppressFBWarnings.java | 38 +++++++++++++++++++
.../flume/conf/internal/package-info.java | 22 +++++++++++
flume-ng-core/pom.xml | 5 ---
.../channel/LoadBalancingChannelSelector.java | 2 +-
.../channel/MultiplexingChannelSelector.java | 2 +-
.../flume/lifecycle/LifecycleSupervisor.java | 2 +-
.../sink/AbstractSingleSinkProcessor.java | 2 +-
.../flume/sink/AbstractSinkProcessor.java | 2 +-
.../flume/sink/FailoverSinkProcessor.java | 2 +-
.../sink/LoadBalancingSinkProcessor.java | 2 +-
.../org/apache/flume/source/ExecSource.java | 2 +-
flume-ng-node/pom.xml | 6 ---
.../node/AbstractConfigurationProvider.java | 2 +-
flume-third-party/pom.xml | 7 ----
15 files changed, 70 insertions(+), 28 deletions(-)
create mode 100644 flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/SuppressFBWarnings.java
create mode 100644 flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/package-info.java
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
index a7bc35fe39..29c40a8283 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/FlumeConfiguration.java
@@ -61,6 +61,7 @@
import org.apache.flume.conf.channel.ChannelType;
import org.apache.flume.conf.configfilter.ConfigFilterConfiguration;
import org.apache.flume.conf.configfilter.ConfigFilterType;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.conf.sink.SinkConfiguration;
import org.apache.flume.conf.sink.SinkGroupConfiguration;
import org.apache.flume.conf.sink.SinkType;
@@ -69,7 +70,6 @@
import org.apache.flume.configfilter.ConfigFilter;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
-import org.apache.logging.log4j.internal.annotation.SuppressFBWarnings;
/**
*
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/SuppressFBWarnings.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/SuppressFBWarnings.java
new file mode 100644
index 0000000000..733eeaf55d
--- /dev/null
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/SuppressFBWarnings.java
@@ -0,0 +1,38 @@
+/*
+ * 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.conf.internal;
+
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+/**
+ * Used to suppress SpotBugs warnings in Flume artifacts.
+ *
+ *
SpotBugs recognizes any annotation whose simple name is {@code SuppressFBWarnings}, so this
+ * type lets us drop the dependency on the LGPL-licensed {@code spotbugs-annotations} artifact.
+ *
+ * This type is not exported via JPMS. Do not use in third-party modules.
+ */
+@Retention(RetentionPolicy.CLASS)
+@Target({ElementType.TYPE, ElementType.FIELD, ElementType.METHOD, ElementType.CONSTRUCTOR, ElementType.PARAMETER})
+public @interface SuppressFBWarnings {
+
+ /** The set of SpotBugs warnings to suppress. */
+ String[] value() default {};
+}
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/package-info.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/package-info.java
new file mode 100644
index 0000000000..c0cd674bd9
--- /dev/null
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/internal/package-info.java
@@ -0,0 +1,22 @@
+/*
+ * 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.
+ */
+/**
+ * Contains types used only by Flume modules.
+ *
+ * These types are not exported via JPMS and are not available to third-party modules.
+ */
+package org.apache.flume.conf.internal;
diff --git a/flume-ng-core/pom.xml b/flume-ng-core/pom.xml
index 7bb547389d..cb8129baab 100644
--- a/flume-ng-core/pom.xml
+++ b/flume-ng-core/pom.xml
@@ -143,11 +143,6 @@
mockito-core
test
-
- com.github.spotbugs
- spotbugs-annotations
- provided
-
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
index cb2b9242ce..78e656a2a6 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/LoadBalancingChannelSelector.java
@@ -18,7 +18,6 @@
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Collections;
import java.util.List;
import java.util.Random;
@@ -26,6 +25,7 @@
import org.apache.flume.Channel;
import org.apache.flume.Context;
import org.apache.flume.Event;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
/**
* Load balancing channel selector. This selector allows for load balancing
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
index 64ea9276b6..d92d4cb05d 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/MultiplexingChannelSelector.java
@@ -16,7 +16,6 @@
*/
package org.apache.flume.channel;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@@ -25,6 +24,7 @@
import org.apache.flume.Context;
import org.apache.flume.Event;
import org.apache.flume.FlumeException;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR")
public class MultiplexingChannelSelector extends AbstractChannelSelector {
diff --git a/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java b/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
index 8151c0c6fd..66b96caa08 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/lifecycle/LifecycleSupervisor.java
@@ -18,7 +18,6 @@
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.HashMap;
import java.util.Map;
import java.util.Map.Entry;
@@ -27,6 +26,7 @@
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.flume.FlumeException;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
index d3f453b3ec..f35a111aff 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSingleSinkProcessor.java
@@ -17,10 +17,10 @@
package org.apache.flume.sink;
import com.google.common.base.Preconditions;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.List;
import org.apache.flume.Sink;
import org.apache.flume.SinkProcessor;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.lifecycle.LifecycleState;
/**
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
index e57cfa008d..a9bbd500c2 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/AbstractSinkProcessor.java
@@ -16,12 +16,12 @@
*/
package org.apache.flume.sink;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import org.apache.flume.Sink;
import org.apache.flume.SinkProcessor;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.lifecycle.LifecycleState;
/**
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
index 49961a30dc..b12f555c98 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/FailoverSinkProcessor.java
@@ -16,7 +16,6 @@
*/
package org.apache.flume.sink;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -29,6 +28,7 @@
import org.apache.flume.EventDeliveryException;
import org.apache.flume.Sink;
import org.apache.flume.Sink.Status;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java b/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
index ad6f0609f8..4ee638c626 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/sink/LoadBalancingSinkProcessor.java
@@ -17,7 +17,6 @@
package org.apache.flume.sink;
import com.google.common.base.Preconditions;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.Iterator;
import java.util.List;
import org.apache.flume.Context;
@@ -26,6 +25,7 @@
import org.apache.flume.Sink;
import org.apache.flume.Sink.Status;
import org.apache.flume.conf.Configurable;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.lifecycle.LifecycleAware;
import org.apache.flume.util.OrderSelector;
import org.apache.flume.util.RandomOrderSelector;
diff --git a/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java b/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
index 2b11df2f06..c99e8dd5f6 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/source/ExecSource.java
@@ -18,7 +18,6 @@
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
@@ -40,6 +39,7 @@
import org.apache.flume.channel.ChannelProcessor;
import org.apache.flume.conf.BatchSizeSupported;
import org.apache.flume.conf.Configurable;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.event.EventBuilder;
import org.apache.flume.instrumentation.SourceCounter;
import org.apache.logging.log4j.LogManager;
diff --git a/flume-ng-node/pom.xml b/flume-ng-node/pom.xml
index 8ac7986bb3..0542d11712 100644
--- a/flume-ng-node/pom.xml
+++ b/flume-ng-node/pom.xml
@@ -107,12 +107,6 @@
provided
-
- com.github.spotbugs
- spotbugs-annotations
- provided
-
-
org.junit.jupiter
junit-jupiter-api
diff --git a/flume-ng-node/src/main/java/org/apache/flume/node/AbstractConfigurationProvider.java b/flume-ng-node/src/main/java/org/apache/flume/node/AbstractConfigurationProvider.java
index 58179eeb2b..cbf958de9c 100644
--- a/flume-ng-node/src/main/java/org/apache/flume/node/AbstractConfigurationProvider.java
+++ b/flume-ng-node/src/main/java/org/apache/flume/node/AbstractConfigurationProvider.java
@@ -21,7 +21,6 @@
import com.google.common.collect.ListMultimap;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
-import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -57,6 +56,7 @@
import org.apache.flume.conf.FlumeConfiguration.AgentConfiguration;
import org.apache.flume.conf.TransactionCapacitySupported;
import org.apache.flume.conf.channel.ChannelSelectorConfiguration;
+import org.apache.flume.conf.internal.SuppressFBWarnings;
import org.apache.flume.conf.sink.SinkConfiguration;
import org.apache.flume.conf.sink.SinkGroupConfiguration;
import org.apache.flume.conf.source.SourceConfiguration;
diff --git a/flume-third-party/pom.xml b/flume-third-party/pom.xml
index 73029cb1d9..b92e7b5b22 100644
--- a/flume-third-party/pom.xml
+++ b/flume-third-party/pom.xml
@@ -85,7 +85,6 @@
1.7.0
2.0.17
1.1.10.8
- 4.9.3
1.19.0
1.1.3
3.8.6
@@ -479,12 +478,6 @@
jsr305
${jsr305.version}
-
-
- com.github.spotbugs
- spotbugs-annotations
- ${spotbugs-annotations.version}
-
From 83794692d3dabd09765abeb7c7ab0fafd037b02f Mon Sep 17 00:00:00 2001
From: "Piotr P. Karwasz"
Date: Fri, 26 Jun 2026 17:35:37 +0200
Subject: [PATCH 5/7] Apply suggestions from code review
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
---
.../java/org/apache/flume/conf/channel/ChannelType.java | 8 --------
1 file changed, 8 deletions(-)
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
index 315c24c5ac..f393314735 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
@@ -64,14 +64,6 @@ private ChannelType(String channelClassName) {
this.channelClassName = channelClassName;
}
- private static ChannelType getChannelTypeByName(String type) {
- for (ChannelType channelType : ChannelType.values()) {
- if (channelType.channelClassName.equals(type)) {
- return channelType;
- }
- }
- return null;
- }
@Deprecated
public String getChannelClassName() {
From 4c658dc1b5b4c4aaf96ec871ba3bf1597a1e6c0d Mon Sep 17 00:00:00 2001
From: Ralph Goers
Date: Fri, 26 Jun 2026 12:20:34 -0700
Subject: [PATCH 6/7] Fix formatting
---
.../src/main/java/org/apache/flume/conf/channel/ChannelType.java | 1 -
1 file changed, 1 deletion(-)
diff --git a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
index f393314735..e6797dc592 100644
--- a/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
+++ b/flume-ng-configuration/src/main/java/org/apache/flume/conf/channel/ChannelType.java
@@ -64,7 +64,6 @@ private ChannelType(String channelClassName) {
this.channelClassName = channelClassName;
}
-
@Deprecated
public String getChannelClassName() {
return channelClassName;
From bc1d39304b50f140406e8409a5a9229e2bdec930 Mon Sep 17 00:00:00 2001
From: Ralph Goers
Date: Fri, 26 Jun 2026 13:44:31 -0700
Subject: [PATCH 7/7] Fix potential NPE
---
.../flume/channel/RoutableProxyChannelSelector.java | 12 ++++++++++--
1 file changed, 10 insertions(+), 2 deletions(-)
diff --git a/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java b/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
index 8fcda6c47c..00369df10d 100644
--- a/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
+++ b/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java
@@ -93,12 +93,20 @@ public void configure(Context context) {
@Override
public List getRequiredChannels(Event event) {
- return getSelector(event).getRequiredChannels(event);
+ ChannelSelector selector = getSelector(event);
+ if (selector != null) {
+ return selector.getRequiredChannels(event);
+ }
+ return new ArrayList<>();
}
@Override
public List getOptionalChannels(Event event) {
- return getSelector(event).getOptionalChannels(event);
+ ChannelSelector selector = getSelector(event);
+ if (selector != null) {
+ return selector.getOptionalChannels(event);
+ }
+ return new ArrayList<>();
}
private ChannelSelector getSelector(Event event) {