From d58f48fd7c69780fee34c6295b3423e5eeb6a74c Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Tue, 18 Jan 2022 22:40:01 +0800 Subject: [PATCH] [pulsar-io] pass client builder if no service url provided to debezium connector (#12145) # Conflicts: # pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/DebeziumSource.java # pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java # tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMySqlSourceTester.java # tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java --- .../pulsar/io/debezium/DebeziumSource.java | 22 +++++++++---------- .../io/debezium/PulsarDatabaseHistory.java | 7 +++--- .../debezium/DebeziumMySqlSourceTester.java | 11 ++++++---- .../debezium/PulsarDebeziumSourcesTest.java | 15 +++++++++---- 4 files changed, 32 insertions(+), 23 deletions(-) diff --git a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/DebeziumSource.java b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/DebeziumSource.java index 0f8d0280a4fec..4f5958ef6e32d 100644 --- a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/DebeziumSource.java +++ b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/DebeziumSource.java @@ -22,6 +22,7 @@ import io.debezium.relational.HistorizedRelationalDatabaseConnectorConfig; import io.debezium.relational.history.DatabaseHistory; +import java.util.Map; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.io.core.SourceContext; @@ -50,10 +51,7 @@ public static void throwExceptionIfConfigNotMatch(Map config, } public static void setConfigIfNull(Map config, String key, String value) { - Object orig = config.get(key); - if (orig == null) { - config.put(key, value); - } + config.putIfAbsent(key, value); } // namespace for output topics, default value is "tenant/namespace" @@ -79,12 +77,8 @@ public void open(Map config, SourceContext sourceContext) throws // database.history : implementation class for database history. setConfigIfNull(config, HistorizedRelationalDatabaseConnectorConfig.DATABASE_HISTORY.name(), DEFAULT_HISTORY); - // database.history.pulsar.service.url, this is set as the value of pulsar.service.url if null. - String serviceUrl = (String) config.get(PulsarKafkaWorkerConfig.PULSAR_SERVICE_URL_CONFIG); - if (serviceUrl == null) { - throw new IllegalArgumentException("Pulsar service URL not provided."); - } - setConfigIfNull(config, PulsarDatabaseHistory.SERVICE_URL.name(), serviceUrl); + // database.history.pulsar.service.url + String pulsarUrl = (String) config.get(PulsarDatabaseHistory.SERVICE_URL.name()); String topicNamespace = topicNamespace(sourceContext); // topic.namespace @@ -98,8 +92,12 @@ public void open(Map config, SourceContext sourceContext) throws setConfigIfNull(config, PulsarKafkaWorkerConfig.OFFSET_STORAGE_TOPIC_CONFIG, topicNamespace + "/" + sourceName + "-" + DEFAULT_OFFSET_TOPIC); - config.put(DatabaseHistory.CONFIGURATION_FIELD_PREFIX_STRING + "pulsar.client.builder", - SerDeUtils.serialize(sourceContext.getPulsarClientBuilder())); + // pass pulsar.client.builder if database.history.pulsar.service.url is not provided + if (StringUtils.isEmpty(pulsarUrl)) { + String pulsarClientBuilder = SerDeUtils.serialize(sourceContext.getPulsarClientBuilder()); + config.put(PulsarDatabaseHistory.CLIENT_BUILDER.name(), pulsarClientBuilder); + } + super.open(config, sourceContext); } diff --git a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java index d8b37c1a1b3b6..ebca3b2db1f11 100644 --- a/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java +++ b/pulsar-io/debezium/core/src/main/java/org/apache/pulsar/io/debezium/PulsarDatabaseHistory.java @@ -104,12 +104,13 @@ public void configure( + getClass().getSimpleName() + "; check the logs for details"); } this.topicName = config.getString(TOPIC); - if (config.getString(CLIENT_BUILDER) == null && config.getString(SERVICE_URL) == null) { + + String clientBuilderBase64Encoded = config.getString(CLIENT_BUILDER); + if (isBlank(clientBuilderBase64Encoded) && isBlank(config.getString(SERVICE_URL))) { throw new IllegalArgumentException("Neither Pulsar Service URL nor ClientBuilder provided."); } - String clientBuilderBase64Encoded = config.getString(CLIENT_BUILDER); this.clientBuilder = PulsarClient.builder(); - if (null != clientBuilderBase64Encoded) { + if (!isBlank(clientBuilderBase64Encoded)) { // deserialize the client builder to the same classloader this.clientBuilder = (ClientBuilder) SerDeUtils.deserialize(clientBuilderBase64Encoded, this.clientBuilder.getClass().getClassLoader()); } else { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMySqlSourceTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMySqlSourceTester.java index ec7e07ccf0be2..dd846341a60b5 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMySqlSourceTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMySqlSourceTester.java @@ -49,7 +49,8 @@ public class DebeziumMySqlSourceTester extends SourceTester