Repository: kafka Updated Branches: refs/heads/0.10.2 1f4f9c74c -> 4f6135931
KAFKA-4700: Don't drop security configs in `StreamsKafkaClient` Author: Ismael Juma <[email protected]> Reviewers: Eno Thereska, Damian Guy, Guozhang Wang Closes #2441 from ijuma/streams-kafka-client-drops-security-configs (cherry picked from commit 31716590454e1bdd53a65836741f316bf4715c9d) Signed-off-by: Guozhang Wang <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/kafka/repo Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/4f613593 Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/4f613593 Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/4f613593 Branch: refs/heads/0.10.2 Commit: 4f61359318f2e7c25ec46d99fd134cf61de28a93 Parents: 1f4f9c7 Author: Ismael Juma <[email protected]> Authored: Thu Jan 26 14:18:54 2017 -0800 Committer: Guozhang Wang <[email protected]> Committed: Thu Jan 26 14:19:02 2017 -0800 ---------------------------------------------------------------------- .../org/apache/kafka/streams/StreamsConfig.java | 7 ++++ .../processor/internals/StreamsKafkaClient.java | 23 ++++++++++- .../internals/StreamsKafkaClientTest.java | 42 ++++++++++++++++++++ 3 files changed, 71 insertions(+), 1 deletion(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java ---------------------------------------------------------------------- diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index 956ae0b..d7d6566 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -388,6 +388,13 @@ public class StreamsConfig extends AbstractConfig { return PRODUCER_PREFIX + producerProp; } + /** + * Returns a copy of the config definition. + */ + public static ConfigDef configDef() { + return new ConfigDef(CONFIG); + } + public StreamsConfig(Map<?, ?> props) { super(CONFIG, props); } http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java ---------------------------------------------------------------------- diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java index df590fb..afde63f 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java @@ -24,6 +24,8 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.Node; +import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.metrics.JmxReporter; import org.apache.kafka.common.metrics.MetricConfig; import org.apache.kafka.common.metrics.Metrics; @@ -52,15 +54,34 @@ import java.util.concurrent.TimeUnit; public class StreamsKafkaClient { + private static final ConfigDef CONFIG = StreamsConfig.configDef() + .withClientSslSupport() + .withClientSaslSupport(); + + public static class Config extends AbstractConfig { + + public static Config fromStreamsConfig(StreamsConfig streamsConfig) { + return new Config(streamsConfig.originals()); + } + + public Config(Map<?, ?> originals) { + super(CONFIG, originals, false); + } + } + private final KafkaClient kafkaClient; private final List<MetricsReporter> reporters; - private final StreamsConfig streamsConfig; + private final Config streamsConfig; private static final int MAX_INFLIGHT_REQUESTS = 100; public StreamsKafkaClient(final StreamsConfig streamsConfig) { + this(Config.fromStreamsConfig(streamsConfig)); + } + public StreamsKafkaClient(final Config streamsConfig) { this.streamsConfig = streamsConfig; + final Time time = new SystemTime(); final Map<String, String> metricTags = new LinkedHashMap<>(); http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java new file mode 100644 index 0000000..2fb5724 --- /dev/null +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java @@ -0,0 +1,42 @@ +/** + * 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.kafka.streams.processor.internals; + +import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.SaslConfigs; +import org.apache.kafka.streams.StreamsConfig; +import org.junit.Test; + +import java.util.Properties; + +import static java.util.Arrays.asList; +import static org.junit.Assert.assertEquals; + +public class StreamsKafkaClientTest { + + @Test + public void testConfigFromStreamsConfig() { + for (final String expectedMechanism : asList("PLAIN", "SCRAM-SHA-512")) { + final Properties props = new Properties(); + props.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, "some_app_id"); + props.setProperty(SaslConfigs.SASL_MECHANISM, expectedMechanism); + props.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9000"); + final StreamsConfig streamsConfig = new StreamsConfig(props); + final AbstractConfig config = StreamsKafkaClient.Config.fromStreamsConfig(streamsConfig); + assertEquals(expectedMechanism, config.values().get(SaslConfigs.SASL_MECHANISM)); + assertEquals(expectedMechanism, config.getString(SaslConfigs.SASL_MECHANISM)); + } + } + +}
