This is an automated email from the ASF dual-hosted git repository. rgoers pushed a commit to branch spring-boot4 in repository https://gitbox.apache.org/repos/asf/logging-flume.git
commit fb9c350c3346f7804fffad91176af8d0cc9269d5 Author: Ralph Goers <[email protected]> AuthorDate: Sat Jun 20 06:39:08 2026 -0700 Add RoutableProxyChannelSelector and allow Interceptors to be configured in Spring Boot --- .../flume/conf/channel/ChannelSelectorType.java | 6 +- flume-ng-core/pom.xml | 5 + .../org/apache/flume/channel/ChannelProcessor.java | 7 ++ .../channel/LoadBalancingChannelSelector.java | 2 + .../flume/channel/MultiplexingChannelSelector.java | 2 + .../channel/RoutableProxyChannelSelector.java | 119 +++++++++++++++++++++ .../flume/lifecycle/LifecycleSupervisor.java | 2 + .../flume/sink/AbstractSingleSinkProcessor.java | 2 + .../apache/flume/sink/AbstractSinkProcessor.java | 2 + .../apache/flume/sink/FailoverSinkProcessor.java | 2 + .../flume/sink/LoadBalancingSinkProcessor.java | 2 + .../java/org/apache/flume/source/ExecSource.java | 2 + .../channel/TestRoutableProxyChannelSelector.java | 111 +++++++++++++++++++ flume-parent/pom.xml | 6 +- 14 files changed, 266 insertions(+), 4 deletions(-) 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 37bc1bcd..c15a8a79 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 a41891d2..c6429bba 100644 --- a/flume-ng-core/pom.xml +++ b/flume-ng-core/pom.xml @@ -144,6 +144,11 @@ <artifactId>mockito-core</artifactId> <scope>test</scope> </dependency> + <dependency> + <groupId>com.github.spotbugs</groupId> + <artifactId>spotbugs-annotations</artifactId> + <scope>provided</scope> + </dependency> </dependencies> <build> 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 c3d85ed4..30ef43b5 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<Interceptor> 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 5aea76d2..cb2b9242 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 @@ package org.apache.flume.channel; 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 @@ import org.apache.flume.Event; * defaults to <tt>ROUND_ROBIN</tt> type, but can be overridden via * configuration.</p> */ +@SuppressFBWarnings("UWF_FIELD_NOT_INITIALIZED_IN_CONSTRUCTOR") public class LoadBalancingChannelSelector extends AbstractChannelSelector { private final List<Channel> 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 2a500980..64ea9276 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.Context; 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 00000000..8fcda6c4 --- /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<String, ChannelSelector> 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<Channel> allChannels = getAllChannels(); + for (Map.Entry<String, String> 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<String, String> map = context.getSubProperties(key + SELECTOR); + if (map != null) { + String channelNames = getRequiredNonNull(map, CHANNELS, key); + Set<String> channelSet = new HashSet<>(Arrays.asList(channelNames.split("\\s+"))); + List<Channel> 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<Channel> getRequiredChannels(Event event) { + return getSelector(event).getRequiredChannels(event); + } + + @Override + public List<Channel> 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<String, String> 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 67ed39ac..8151c0c6 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 @@ package org.apache.flume.lifecycle; 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.flume.FlumeException; 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 343ee73b..d3f453b3 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 @@ import org.apache.flume.lifecycle.LifecycleState; /** * 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 2be6b160..e57cfa00 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 @@ import org.apache.flume.lifecycle.LifecycleState; /** * 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 a09a3107..49961a30 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 @@ import org.apache.logging.log4j.Logger; * 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 30943ca4..ad6f0609 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 @@ import org.apache.logging.log4j.Logger; * @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 3a6778a2..2b11df2f 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 @@ package org.apache.flume.source; 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 @@ import org.apache.logging.log4j.Logger; * TODO * </p> */ +@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 00000000..f1df71cb --- /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<Channel> channels = new ArrayList<Channel>(); + 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<Channel> 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<String, String> 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<Channel> channels1 = new ArrayList<>(); + channels1.add(channels.get(0)); + channels1.add(channels.get(1)); + LoadBalancingChannelSelector loadBalancingSelector = + (LoadBalancingChannelSelector) ChannelSelectorFactory.create(channels1, config); + selector.setDefaultSelector(loadBalancingSelector); + List<Channel> 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<Channel> 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<String, String> 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 261737b6..57bc9d9c 100644 --- a/flume-parent/pom.xml +++ b/flume-parent/pom.xml @@ -237,9 +237,9 @@ <commons-text.version>1.15.0</commons-text.version> <curator.version>5.9.0</curator.version> <derby.version>10.17.1.0</derby.version> - <jackson-annotations.version>2.19.0</jackson-annotations.version> - <fasterxml.jackson.version>2.21.1</fasterxml.jackson.version> - <fasterxml.jackson.databind.version>2.21.1</fasterxml.jackson.databind.version> + <jackson-annotations.version>2.22</jackson-annotations.version> + <fasterxml.jackson.version>2.22.0</fasterxml.jackson.version> + <fasterxml.jackson.databind.version>2.22.0</fasterxml.jackson.databind.version> <fest-reflect.version>1.4.1</fest-reflect.version> <gson.version>2.14.0</gson.version> <guava.version>33.4.8-jre</guava.version>
