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) {