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..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; @@ -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/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-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..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 @@ -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 * 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 a41891d265..cb8129baab 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-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..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 @@ -25,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 @@ -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..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 @@ -24,7 +24,9 @@ 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 { 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..00369df10d --- /dev/null +++ b/flume-ng-core/src/main/java/org/apache/flume/channel/RoutableProxyChannelSelector.java @@ -0,0 +1,127 @@ +/* + * 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) { + ChannelSelector selector = getSelector(event); + if (selector != null) { + return selector.getRequiredChannels(event); + } + return new ArrayList<>(); + } + + @Override + public List getOptionalChannels(Event event) { + ChannelSelector selector = getSelector(event); + if (selector != null) { + return selector.getOptionalChannels(event); + } + return new ArrayList<>(); + } + + 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..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 @@ -26,9 +26,11 @@ 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; +@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..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 @@ -20,11 +20,13 @@ 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; /** * 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..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 @@ -21,11 +21,13 @@ 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; /** * 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..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 @@ -28,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; @@ -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..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 @@ -25,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; @@ -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..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 @@ -39,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; @@ -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-ng-node/pom.xml b/flume-ng-node/pom.xml index c746cc4f68..0542d11712 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,13 +104,6 @@ com.google.code.findbugs jsr305 - ${jsr305.version} - provided - - - - com.github.spotbugs - spotbugs-annotations provided 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-parent/pom.xml b/flume-parent/pom.xml index 261737b658..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.19.0 - 2.21.1 - 2.21.1 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..b92e7b5b22 --- /dev/null +++ b/flume-third-party/pom.xml @@ -0,0 +1,520 @@ + + + + + 4.0.0 + + + org.apache.flume + flume-project + 2.0.0-SNAPSHOT + ../pom.xml + + 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 + 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} + + + + + + + + + + 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