This is an automated email from the ASF dual-hosted git repository.
sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 8599151 Issue 1943: remove serializable from builders and add
loadData to load configuration from a config map (#2004)
8599151 is described below
commit 8599151eeed3bcec78964f40b21793220e0b27ea
Author: Sijie Guo <[email protected]>
AuthorDate: Thu Jun 21 12:23:41 2018 -0700
Issue 1943: remove serializable from builders and add loadData to load
configuration from a config map (#2004)
* Issue 1943: remove serializable from builders and add loadData to load
configuration from a config map
*Motivation*
Make pulsar users be able to use java Serialization framework for
distributing client settings. This is required for using pulsar in frameworks
like Spark.
*Changes*
- Remove serializable from builders
- Add a tool to load client/producer/consumer/reader configuration data
from a Map config (which is serializable)
Fixes #1943
* Address review comments:
- add example in javadoc
- change client builder to `loadConf` as well
- change ConfigurationDataUtil to merge config to existing settings
Signed-off-by: Sijie Guo <[email protected]>
* set FAIL_ON_UNKNOWN_PROPERTIES to true
add a test case to make sure exceptions thrown when loading config with
unknown properties
---
.../apache/pulsar/client/api/ClientBuilder.java | 23 ++++-
.../apache/pulsar/client/api/ConsumerBuilder.java | 23 ++++-
.../pulsar/client/api/ConsumerEventListener.java | 4 +-
.../apache/pulsar/client/api/CryptoKeyReader.java | 3 +-
.../org/apache/pulsar/client/api/MessageId.java | 3 +-
.../apache/pulsar/client/api/ProducerBuilder.java | 23 ++++-
.../org/apache/pulsar/client/api/PulsarClient.java | 3 +
.../apache/pulsar/client/api/ReaderBuilder.java | 24 ++++-
.../pulsar/client/impl/ClientBuilderImpl.java | 14 ++-
.../pulsar/client/impl/ConsumerBuilderImpl.java | 13 ++-
.../pulsar/client/impl/ProducerBuilderImpl.java | 12 ++-
.../pulsar/client/impl/ReaderBuilderImpl.java | 12 ++-
.../client/impl/conf/ClientConfigurationData.java | 2 +
.../client/impl/conf/ConfigurationDataUtils.java | 75 ++++++++++++++
.../impl/conf/ProducerConfigurationData.java | 1 +
.../impl/conf/ConfigurationDataUtilsTest.java | 111 +++++++++++++++++++++
.../MultiConsumersOneOutputTopicProducersTest.java | 5 +
17 files changed, 326 insertions(+), 25 deletions(-)
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
index 071e292..e1fe3ec 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
@@ -18,7 +18,6 @@
*/
package org.apache.pulsar.client.api;
-import java.io.Serializable;
import java.util.Map;
import java.util.concurrent.TimeUnit;
@@ -29,7 +28,7 @@ import
org.apache.pulsar.client.api.PulsarClientException.UnsupportedAuthenticat
*
* @since 2.0.0
*/
-public interface ClientBuilder extends Serializable, Cloneable {
+public interface ClientBuilder extends Cloneable {
/**
* @return the new {@link PulsarClient} instance
@@ -37,6 +36,26 @@ public interface ClientBuilder extends Serializable,
Cloneable {
PulsarClient build() throws PulsarClientException;
/**
+ * Load the configuration from provided <tt>config</tt> map.
+ *
+ * <p>Example:
+ * <pre>
+ * Map<String, Object> config = new HashMap<>();
+ * config.put("serviceUrl", "pulsar://localhost:5550");
+ * config.put("numIoThreads", 20);
+ *
+ * ClientBuilder builder = ...;
+ * builder = builder.loadConf(config);
+ *
+ * PulsarClient client = builder.build();
+ * </pre>
+ *
+ * @param config configuration to load
+ * @return client builder instance
+ */
+ ClientBuilder loadConf(Map<String, Object> config);
+
+ /**
* Create a copy of the current client builder.
* <p>
* Cloning the builder can be used to share an incomplete configuration
and specialize it multiple times. For
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java
index f0014a5..8657859 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java
@@ -18,7 +18,6 @@
*/
package org.apache.pulsar.client.api;
-import java.io.Serializable;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
@@ -32,7 +31,7 @@ import java.util.regex.Pattern;
*
* @since 2.0.0
*/
-public interface ConsumerBuilder<T> extends Serializable, Cloneable {
+public interface ConsumerBuilder<T> extends Cloneable {
/**
* Create a copy of the current consumer builder.
@@ -53,6 +52,26 @@ public interface ConsumerBuilder<T> extends Serializable,
Cloneable {
ConsumerBuilder<T> clone();
/**
+ * Load the configuration from provided <tt>config</tt> map.
+ *
+ * <p>Example:
+ * <pre>
+ * Map<String, Object> config = new HashMap<>();
+ * config.put("ackTimeoutMillis", 1000);
+ * config.put("receiverQueueSize", 2000);
+ *
+ * ConsumerBuilder<byte[]> builder = ...;
+ * builder = builder.loadConf(config);
+ *
+ * Consumer<byte[]> consumer = builder.subscribe();
+ * </pre>
+ *
+ * @param config configuration to load
+ * @return consumer builder instance
+ */
+ ConsumerBuilder<T> loadConf(Map<String, Object> config);
+
+ /**
* Finalize the {@link Consumer} creation by subscribing to the topic.
*
* <p>
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerEventListener.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerEventListener.java
index e2e6274..af0bd50 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerEventListener.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ConsumerEventListener.java
@@ -18,10 +18,12 @@
*/
package org.apache.pulsar.client.api;
+import java.io.Serializable;
+
/**
* Listener on the consumer state changes.
*/
-public interface ConsumerEventListener {
+public interface ConsumerEventListener extends Serializable {
/**
* Notified when the consumer group is changed, and the consumer becomes
the active consumer.
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java
index 4acb6e9..274191f 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java
@@ -18,9 +18,10 @@
*/
package org.apache.pulsar.client.api;
+import java.io.Serializable;
import java.util.Map;
-public interface CryptoKeyReader {
+public interface CryptoKeyReader extends Serializable {
/**
* Return the encryption key corresponding to the key name in the argument
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/MessageId.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/MessageId.java
index f74c169..1bb3e08 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/api/MessageId.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/api/MessageId.java
@@ -20,6 +20,7 @@ package org.apache.pulsar.client.api;
import java.io.IOException;
+import java.io.Serializable;
import org.apache.pulsar.client.impl.MessageIdImpl;
/**
@@ -30,7 +31,7 @@ import org.apache.pulsar.client.impl.MessageIdImpl;
*
*
*/
-public interface MessageId extends Comparable<MessageId>{
+public interface MessageId extends Comparable<MessageId>, Serializable {
/**
* Serialize the message ID into a byte array
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java
index a9723a0..8256b4a 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java
@@ -18,7 +18,6 @@
*/
package org.apache.pulsar.client.api;
-import java.io.Serializable;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -30,7 +29,7 @@ import
org.apache.pulsar.client.api.PulsarClientException.ProducerQueueIsFullErr
*
* @see PulsarClient#newProducer()
*/
-public interface ProducerBuilder<T> extends Serializable, Cloneable {
+public interface ProducerBuilder<T> extends Cloneable {
/**
* Finalize the creation of the {@link Producer} instance.
@@ -59,6 +58,26 @@ public interface ProducerBuilder<T> extends Serializable,
Cloneable {
CompletableFuture<Producer<T>> createAsync();
/**
+ * Load the configuration from provided <tt>config</tt> map.
+ *
+ * <p>Example:
+ * <pre>
+ * Map<String, Object> config = new HashMap<>();
+ * config.put("producerName", "test-producer");
+ * config.put("sendTimeoutMs", 2000);
+ *
+ * ProducerBuilder<byte[]> builder = ...;
+ * builder = builder.loadConf(config);
+ *
+ * Producer<byte[]> producer = builder.create();
+ * </pre>
+ *
+ * @param config configuration to load
+ * @return producer builder instance
+ */
+ ProducerBuilder<T> loadConf(Map<String, Object> config);
+
+ /**
* Create a copy of the current {@link ProducerBuilder}.
* <p>
* Cloning the builder can be used to share an incomplete configuration
and specialize it multiple times. For
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/PulsarClient.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/PulsarClient.java
index 98af0b3..5827573 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/api/PulsarClient.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/api/PulsarClient.java
@@ -19,10 +19,13 @@
package org.apache.pulsar.client.api;
import java.io.Closeable;
+import java.util.Map;
import java.util.concurrent.CompletableFuture;
import org.apache.pulsar.client.impl.ClientBuilderImpl;
import org.apache.pulsar.client.impl.PulsarClientImpl;
+import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
+import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils;
/**
* Class that provides a client interface to Pulsar.
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java
index 314af7b..bc162bd 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/ReaderBuilder.java
@@ -18,7 +18,7 @@
*/
package org.apache.pulsar.client.api;
-import java.io.Serializable;
+import java.util.Map;
import java.util.concurrent.CompletableFuture;
/**
@@ -28,7 +28,7 @@ import java.util.concurrent.CompletableFuture;
*
* @since 2.0.0
*/
-public interface ReaderBuilder<T> extends Serializable, Cloneable {
+public interface ReaderBuilder<T> extends Cloneable {
/**
* Finalize the creation of the {@link Reader} instance.
@@ -55,6 +55,26 @@ public interface ReaderBuilder<T> extends Serializable,
Cloneable {
CompletableFuture<Reader<T>> createAsync();
/**
+ * Load the configuration from provided <tt>config</tt> map.
+ *
+ * <p>Example:
+ * <pre>
+ * Map<String, Object> config = new HashMap<>();
+ * config.put("topicName", "test-topic");
+ * config.put("receiverQueueSize", 2000);
+ *
+ * ReaderBuilder<byte[]> builder = ...;
+ * builder = builder.loadConf(config);
+ *
+ * Reader<byte[]> reader = builder.create();
+ * </pre>
+ *
+ * @param config configuration to load
+ * @return reader builder instance
+ */
+ ReaderBuilder<T> loadConf(Map<String, Object> config);
+
+ /**
* Create a copy of the current {@link ReaderBuilder}.
* <p>
* Cloning the builder can be used to share an incomplete configuration
and specialize it multiple times. For
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
index ad4b510..657efd3 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
@@ -28,18 +28,17 @@ import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
import
org.apache.pulsar.client.api.PulsarClientException.UnsupportedAuthenticationException;
import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
+import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils;
public class ClientBuilderImpl implements ClientBuilder {
- private static final long serialVersionUID = 1L;
-
- final ClientConfigurationData conf;
+ ClientConfigurationData conf;
public ClientBuilderImpl() {
this(new ClientConfigurationData());
}
- private ClientBuilderImpl(ClientConfigurationData conf) {
+ public ClientBuilderImpl(ClientConfigurationData conf) {
this.conf = conf;
}
@@ -58,6 +57,13 @@ public class ClientBuilderImpl implements ClientBuilder {
}
@Override
+ public ClientBuilder loadConf(Map<String, Object> config) {
+ conf = ConfigurationDataUtils.loadData(
+ config, conf, ClientConfigurationData.class);
+ return this;
+ }
+
+ @Override
public ClientBuilder serviceUrl(String serviceUrl) {
conf.setServiceUrl(serviceUrl);
return this;
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java
index af11ef5..f0067f7 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java
@@ -39,6 +39,7 @@ import
org.apache.pulsar.client.api.PulsarClientException.InvalidConfigurationEx
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.SubscriptionInitialPosition;
import org.apache.pulsar.client.api.SubscriptionType;
+import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils;
import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData;
import org.apache.pulsar.common.util.FutureUtil;
@@ -48,10 +49,8 @@ import lombok.NonNull;
public class ConsumerBuilderImpl<T> implements ConsumerBuilder<T> {
- private static final long serialVersionUID = 1L;
-
private final PulsarClientImpl client;
- private final ConsumerConfigurationData<T> conf;
+ private ConsumerConfigurationData<T> conf;
private final Schema<T> schema;
private static long MIN_ACK_TIMEOUT_MILLIS = 1000;
@@ -60,13 +59,19 @@ public class ConsumerBuilderImpl<T> implements
ConsumerBuilder<T> {
this(client, new ConsumerConfigurationData<T>(), schema);
}
- private ConsumerBuilderImpl(PulsarClientImpl client,
ConsumerConfigurationData<T> conf, Schema<T> schema) {
+ ConsumerBuilderImpl(PulsarClientImpl client, ConsumerConfigurationData<T>
conf, Schema<T> schema) {
this.client = client;
this.conf = conf;
this.schema = schema;
}
@Override
+ public ConsumerBuilder<T> loadConf(Map<String, Object> config) {
+ this.conf = ConfigurationDataUtils.loadData(config, conf,
ConsumerConfigurationData.class);
+ return this;
+ }
+
+ @Override
public ConsumerBuilder<T> clone() {
return new ConsumerBuilderImpl<>(client, conf.clone(), schema);
}
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java
index 0f5eab0..199665b 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBuilderImpl.java
@@ -33,6 +33,7 @@ import org.apache.pulsar.client.api.ProducerBuilder;
import org.apache.pulsar.client.api.ProducerCryptoFailureAction;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.Schema;
+import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils;
import org.apache.pulsar.client.impl.conf.ProducerConfigurationData;
import org.apache.pulsar.common.util.FutureUtil;
@@ -40,10 +41,8 @@ import lombok.NonNull;
public class ProducerBuilderImpl<T> implements ProducerBuilder<T> {
- private static final long serialVersionUID = 1L;
-
private final PulsarClientImpl client;
- private final ProducerConfigurationData conf;
+ private ProducerConfigurationData conf;
private final Schema<T> schema;
ProducerBuilderImpl(PulsarClientImpl client, Schema<T> schema) {
@@ -89,6 +88,13 @@ public class ProducerBuilderImpl<T> implements
ProducerBuilder<T> {
}
@Override
+ public ProducerBuilder<T> loadConf(Map<String, Object> config) {
+ conf = ConfigurationDataUtils.loadData(
+ config, conf, ProducerConfigurationData.class);
+ return this;
+ }
+
+ @Override
public ProducerBuilder<T> topic(String topicName) {
conf.setTopicName(topicName);
return this;
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderBuilderImpl.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderBuilderImpl.java
index 9e8bce8..d74dc83 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderBuilderImpl.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderBuilderImpl.java
@@ -18,6 +18,7 @@
*/
package org.apache.pulsar.client.impl;
+import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
@@ -29,16 +30,15 @@ import org.apache.pulsar.client.api.Reader;
import org.apache.pulsar.client.api.ReaderBuilder;
import org.apache.pulsar.client.api.ReaderListener;
import org.apache.pulsar.client.api.Schema;
+import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils;
import org.apache.pulsar.client.impl.conf.ReaderConfigurationData;
import org.apache.pulsar.common.util.FutureUtil;
public class ReaderBuilderImpl<T> implements ReaderBuilder<T> {
- private static final long serialVersionUID = 1L;
-
private final PulsarClientImpl client;
- private final ReaderConfigurationData<T> conf;
+ private ReaderConfigurationData<T> conf;
private final Schema<T> schema;
@@ -95,6 +95,12 @@ public class ReaderBuilderImpl<T> implements
ReaderBuilder<T> {
}
@Override
+ public ReaderBuilder<T> loadConf(Map<String, Object> config) {
+ conf = ConfigurationDataUtils.loadData(config, conf,
ReaderConfigurationData.class);
+ return this;
+ }
+
+ @Override
public ReaderBuilder<T> topic(String topicName) {
conf.setTopicName(topicName);
return this;
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
index cedc892..b51d61f 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
@@ -63,4 +63,6 @@ public class ClientConfigurationData implements Serializable,
Cloneable {
throw new RuntimeException("Failed to clone
ClientConfigurationData");
}
}
+
+
}
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtils.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtils.java
new file mode 100644
index 0000000..4523939
--- /dev/null
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtils.java
@@ -0,0 +1,75 @@
+/**
+ * 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.pulsar.client.impl.conf;
+
+import com.fasterxml.jackson.annotation.JsonInclude.Include;
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.google.common.collect.Maps;
+import io.netty.util.concurrent.FastThreadLocal;
+import java.io.IOException;
+import java.util.Map;
+
+/**
+ * Utils for loading configuration data.
+ */
+public final class ConfigurationDataUtils {
+
+ public static ObjectMapper create() {
+ ObjectMapper mapper = new ObjectMapper();
+ // forward compatibility for the properties may go away in the future
+ mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES,
true);
+
mapper.configure(DeserializationFeature.READ_UNKNOWN_ENUM_VALUES_AS_NULL,
false);
+ mapper.setSerializationInclusion(Include.NON_NULL);
+ return mapper;
+ }
+
+ private static final FastThreadLocal<ObjectMapper> mapper = new
FastThreadLocal<ObjectMapper>() {
+ @Override
+ protected ObjectMapper initialValue() throws Exception {
+ return create();
+ }
+ };
+
+ public static ObjectMapper getThreadLocal() {
+ return mapper.get();
+ }
+
+ private ConfigurationDataUtils() {}
+
+ public static <T> T loadData(Map<String, Object> config,
+ T existingData,
+ Class<T> dataCls) {
+ ObjectMapper mapper = getThreadLocal();
+ try {
+ String existingConfigJson =
mapper.writeValueAsString(existingData);
+ Map<String, Object> existingConfig =
mapper.readValue(existingConfigJson, Map.class);
+ Map<String, Object> newConfig = Maps.newHashMap();
+ newConfig.putAll(existingConfig);
+ newConfig.putAll(config);
+ String configJson = mapper.writeValueAsString(newConfig);
+ return mapper.readValue(configJson, dataCls);
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to load config into existing
configuration data", e);
+ }
+
+ }
+
+}
+
diff --git
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java
index 9f06809..57ca135 100644
---
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java
+++
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ProducerConfigurationData.java
@@ -80,6 +80,7 @@ public class ProducerConfigurationData implements
Serializable, Cloneable {
* Returns true if encryption keys are added
*
*/
+ @JsonIgnore
public boolean isEncryptionEnabled() {
return (this.encryptionKeys != null) && !this.encryptionKeys.isEmpty()
&& (this.cryptoKeyReader != null);
}
diff --git
a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtilsTest.java
b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtilsTest.java
new file mode 100644
index 0000000..83ed4e8
--- /dev/null
+++
b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/conf/ConfigurationDataUtilsTest.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.pulsar.client.impl.conf;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.fail;
+
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
+import org.testng.annotations.Test;
+
+/**
+ * Unit test {@link ConfigurationDataUtils}.
+ */
+public class ConfigurationDataUtilsTest {
+
+ @Test
+ public void testLoadClientConfigurationData() {
+ ClientConfigurationData confData = new ClientConfigurationData();
+ confData.setServiceUrl("pulsar://unknown:6650");
+ confData.setMaxLookupRequest(600);
+ confData.setNumIoThreads(33);
+ Map<String, Object> config = new HashMap<>();
+ config.put("serviceUrl", "pulsar://localhost:6650");
+ config.put("maxLookupRequest", 70000);
+ confData = ConfigurationDataUtils.loadData(config, confData,
ClientConfigurationData.class);
+ assertEquals("pulsar://localhost:6650", confData.getServiceUrl());
+ assertEquals(70000, confData.getMaxLookupRequest());
+ assertEquals(33, confData.getNumIoThreads());
+ }
+
+ @Test
+ public void testLoadProducerConfigurationData() {
+ ProducerConfigurationData confData = new ProducerConfigurationData();
+ confData.setProducerName("unset");
+ confData.setBatchingEnabled(true);
+ confData.setBatchingMaxMessages(1234);
+ Map<String, Object> config = new HashMap<>();
+ config.put("producerName", "test-producer");
+ config.put("batchingEnabled", false);
+ confData = ConfigurationDataUtils.loadData(config, confData,
ProducerConfigurationData.class);
+ assertEquals("test-producer", confData.getProducerName());
+ assertEquals(false, confData.isBatchingEnabled());
+ assertEquals(1234, confData.getBatchingMaxMessages());
+ }
+
+ @Test
+ public void testLoadConsumerConfigurationData() {
+ ConsumerConfigurationData confData = new ConsumerConfigurationData();
+ confData.setSubscriptionName("unknown-subscription");
+ confData.setPriorityLevel(10000);
+ confData.setConsumerName("unknown-consumer");
+ Map<String, Object> config = new HashMap<>();
+ config.put("subscriptionName", "test-subscription");
+ config.put("priorityLevel", 100);
+ confData = ConfigurationDataUtils.loadData(config, confData,
ConsumerConfigurationData.class);
+ assertEquals("test-subscription", confData.getSubscriptionName());
+ assertEquals(100, confData.getPriorityLevel());
+ assertEquals("unknown-consumer", confData.getConsumerName());
+ }
+
+ @Test
+ public void testLoadReaderConfigurationData() {
+ ReaderConfigurationData confData = new ReaderConfigurationData();
+ confData.setTopicName("unknown");
+ confData.setReceiverQueueSize(1000000);
+ confData.setReaderName("unknown-reader");
+ Map<String, Object> config = new HashMap<>();
+ config.put("topicName", "test-topic");
+ config.put("receiverQueueSize", 100);
+ confData = ConfigurationDataUtils.loadData(config, confData,
ReaderConfigurationData.class);
+ assertEquals("test-topic", confData.getTopicName());
+ assertEquals(100, confData.getReceiverQueueSize());
+ assertEquals("unknown-reader", confData.getReaderName());
+ }
+
+ @Test
+ public void testLoadConfigurationDataWithUnknownFields() {
+ ReaderConfigurationData confData = new ReaderConfigurationData();
+ confData.setTopicName("unknown");
+ confData.setReceiverQueueSize(1000000);
+ confData.setReaderName("unknown-reader");
+ Map<String, Object> config = new HashMap<>();
+ config.put("unknown", "test-topic");
+ config.put("receiverQueueSize", 100);
+ try {
+ ConfigurationDataUtils.loadData(config, confData,
ReaderConfigurationData.class);
+ fail("Should fail loading configuration data with unknown fields");
+ } catch (RuntimeException re) {
+ assertTrue(re.getCause() instanceof IOException);
+ }
+ }
+}
diff --git
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/producers/MultiConsumersOneOutputTopicProducersTest.java
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/producers/MultiConsumersOneOutputTopicProducersTest.java
index 81cbbda..def0926 100644
---
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/producers/MultiConsumersOneOutputTopicProducersTest.java
+++
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/producers/MultiConsumersOneOutputTopicProducersTest.java
@@ -89,6 +89,11 @@ public class MultiConsumersOneOutputTopicProducersTest {
}
@Override
+ public ProducerBuilder<byte[]> loadConf(Map<String, Object> config) {
+ return this;
+ }
+
+ @Override
public ProducerBuilder<byte[]> clone() {
return this;
}