From 83dd8e7b1eb23c499d51e72c15d236e006716c20 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Thu, 23 Sep 2021 11:04:23 +0800 Subject: [PATCH 1/6] stash --- .../org/apache/pulsar/io/debezium/DebeziumSource.java | 10 +++++----- .../pulsar/io/debezium/PulsarDatabaseHistory.java | 7 ++++--- 2 files changed, 9 insertions(+), 8 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 b9074b91bc7c5..6f75233fa990b 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 @@ -81,9 +81,6 @@ public void open(Map config, SourceContext sourceContext) throws // database.history.pulsar.service.url String pulsarUrl = (String) config.get(PulsarDatabaseHistory.SERVICE_URL.name()); - if (StringUtils.isEmpty(pulsarUrl)) { - throw new IllegalArgumentException("Pulsar service URL for History Database not provided."); - } String topicNamespace = topicNamespace(sourceContext); // topic.namespace @@ -97,8 +94,11 @@ 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 be152a6da8eb2..a98e817b4f164 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 @@ -103,14 +103,15 @@ public void configure( } this.topicName = config.getString(TOPIC); - if (config.getString(CLIENT_BUILDER) == null && config.getString(SERVICE_URL) == null) { + if (StringUtils.isEmpty(config.getString(CLIENT_BUILDER)) && StringUtils.isEmpty(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 (!StringUtils.isEmpty(clientBuilderBase64Encoded)) { // deserialize the client builder to the same classloader - this.clientBuilder = (ClientBuilder) SerDeUtils.deserialize(clientBuilderBase64Encoded, this.clientBuilder.getClass().getClassLoader()); + this.clientBuilder = (ClientBuilder) SerDeUtils.deserialize(clientBuilderBase64Encoded, + this.clientBuilder.getClass().getClassLoader()); } else { this.clientBuilder.serviceUrl(config.getString(SERVICE_URL)); } From 2e8838b6cd4c9ca49935a3b0ceb1d87d3c71f02d Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Thu, 23 Sep 2021 11:24:39 +0800 Subject: [PATCH 2/6] only pass client builder if no service url provided --- .../org/apache/pulsar/io/debezium/DebeziumSource.java | 9 ++------- .../apache/pulsar/io/debezium/PulsarDatabaseHistory.java | 6 +++--- 2 files changed, 5 insertions(+), 10 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 6f75233fa990b..7e3a0a3b6895e 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 @@ -18,10 +18,8 @@ */ package org.apache.pulsar.io.debezium; -import io.debezium.relational.history.DatabaseHistory; -import java.util.Map; - import io.debezium.relational.HistorizedRelationalDatabaseConnectorConfig; +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 +48,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" 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 a98e817b4f164..7a4812c537496 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 @@ -103,12 +103,12 @@ public void configure( } this.topicName = config.getString(TOPIC); - if (StringUtils.isEmpty(config.getString(CLIENT_BUILDER)) && StringUtils.isEmpty(config.getString(SERVICE_URL))) { + 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 (!StringUtils.isEmpty(clientBuilderBase64Encoded)) { + if (!isBlank(clientBuilderBase64Encoded)) { // deserialize the client builder to the same classloader this.clientBuilder = (ClientBuilder) SerDeUtils.deserialize(clientBuilderBase64Encoded, this.clientBuilder.getClass().getClassLoader()); From 8acf91913ad97d2e1551ac4d3ae4535a3fcc30ed Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Wed, 17 Nov 2021 13:38:25 +0800 Subject: [PATCH 3/6] add test cases in integration tests --- .../debezium/DebeziumMongoDbSourceTester.java | 6 ++- .../debezium/DebeziumMsSqlSourceTester.java | 6 ++- .../debezium/DebeziumMySqlSourceTester.java | 7 ++- .../DebeziumPostgreSqlSourceTester.java | 6 ++- .../debezium/PulsarDebeziumSourcesTest.java | 51 ++++++++++++++----- 5 files changed, 55 insertions(+), 21 deletions(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java index 110ff11c00d07..24eb491348f8b 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java @@ -39,7 +39,7 @@ public class DebeziumMongoDbSourceTester extends SourceTester confirmedFlushLsn = new AtomicReference<>("not read yet"); - public DebeziumPostgreSqlSourceTester(PulsarCluster cluster) { + public DebeziumPostgreSqlSourceTester(PulsarCluster cluster, boolean testWithClientBuilder) { super(NAME); this.pulsarCluster = cluster; /* @@ -80,7 +80,9 @@ public DebeziumPostgreSqlSourceTester(PulsarCluster cluster) { sourceConfig.put("database.dbname", "postgres"); sourceConfig.put("schema.whitelist", "inventory"); sourceConfig.put("table.blacklist", "inventory.spatial_ref_sys,inventory.geom"); - sourceConfig.put("database.history.pulsar.service.url", pulsarServiceUrl); + if (!testWithClientBuilder) { + sourceConfig.put("database.history.pulsar.service.url", pulsarServiceUrl); + } sourceConfig.put("topic.namespace", "debezium/postgresql"); } diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java index 98363195b7be3..319cada12a74a 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java @@ -45,31 +45,52 @@ public class PulsarDebeziumSourcesTest extends PulsarIOTestBase { @Test(groups = "source") public void testDebeziumMySqlSourceJson() throws Exception { - testDebeziumMySqlConnect("org.apache.kafka.connect.json.JsonConverter", true); + testDebeziumMySqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + } + + @Test(groups = "source") + public void testDebeziumMySqlSourceJsonWithClientBuilder() throws Exception { + testDebeziumMySqlConnect("org.apache.kafka.connect.json.JsonConverter", true, true); } @Test(groups = "source") public void testDebeziumMySqlSourceAvro() throws Exception { testDebeziumMySqlConnect( - "org.apache.pulsar.kafka.shade.io.confluent.connect.avro.AvroConverter", false); + "org.apache.pulsar.kafka.shade.io.confluent.connect.avro.AvroConverter", false, false); } @Test(groups = "source") public void testDebeziumPostgreSqlSource() throws Exception { - testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true); + testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + } + + @Test(groups = "source") + public void testDebeziumPostgreSqlSourceWithClientBuilder() throws Exception { + testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, true); } @Test(groups = "source") public void testDebeziumMongoDbSource() throws Exception{ - testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true); + testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + } + + @Test(groups = "source") + public void testDebeziumMongoDbSourceWithClientBuilder() throws Exception{ + testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true, true); } @Test(groups = "source") public void testDebeziumMsSqlSource() throws Exception{ - testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true); + testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + } + + @Test(groups = "source") + public void testDebeziumMsSqlSourceWithClientBuilder() throws Exception{ + testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, falsetrue; } - private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { + private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWithEnvelope, + boolean testWithClientBuilder) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -104,7 +125,7 @@ private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWit admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumMySqlSourceTester sourceTester = new DebeziumMySqlSourceTester(pulsarCluster, converterClassName); + DebeziumMySqlSourceTester sourceTester = new DebeziumMySqlSourceTester(pulsarCluster, converterClassName, testWithClientBuilder); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); // setup debezium mysql server @@ -118,7 +139,8 @@ private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWit runner.testSource(sourceTester); } - private void testDebeziumPostgreSqlConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { + private void testDebeziumPostgreSqlConnect(String converterClassName, boolean jsonWithEnvelope, + boolean testWithClientBuilder) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -142,7 +164,7 @@ private void testDebeziumPostgreSqlConnect(String converterClassName, boolean js admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumPostgreSqlSourceTester sourceTester = new DebeziumPostgreSqlSourceTester(pulsarCluster); + DebeziumPostgreSqlSourceTester sourceTester = new DebeziumPostgreSqlSourceTester(pulsarCluster, testWithClientBuilder); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); // setup debezium postgresql server @@ -156,7 +178,8 @@ private void testDebeziumPostgreSqlConnect(String converterClassName, boolean js runner.testSource(sourceTester); } - private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { + private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonWithEnvelope, + boolean testWithClientBuilder) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -181,7 +204,8 @@ private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonW admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumMongoDbSourceTester sourceTester = new DebeziumMongoDbSourceTester(pulsarCluster); + DebeziumMongoDbSourceTester sourceTester = + new DebeziumMongoDbSourceTester(pulsarCluster, testWithClientBuilder); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); // setup debezium mongodb server @@ -195,7 +219,8 @@ private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonW runner.testSource(sourceTester); } - private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { + private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWithEnvelope, + boolean testWithClientBuilder) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -218,7 +243,7 @@ private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWit admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumMsSqlSourceTester sourceTester = new DebeziumMsSqlSourceTester(pulsarCluster); + DebeziumMsSqlSourceTester sourceTester = new DebeziumMsSqlSourceTester(pulsarCluster, testWithClientBuilder); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); DebeziumMsSqlContainer msSqlContainer = new DebeziumMsSqlContainer(pulsarCluster.getClusterName()); From 4a27ce7854dd2c11614c14fc023518bda5e4e106 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Wed, 17 Nov 2021 14:27:34 +0800 Subject: [PATCH 4/6] fix --- .../io/sources/debezium/PulsarDebeziumSourcesTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java index 319cada12a74a..e5236ee16e522 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java @@ -86,7 +86,7 @@ public void testDebeziumMsSqlSource() throws Exception{ @Test(groups = "source") public void testDebeziumMsSqlSourceWithClientBuilder() throws Exception{ - testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, falsetrue; + testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); } private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWithEnvelope, From 33407ceb8d69ca53760939f7fbd98ec253dbcb07 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Thu, 18 Nov 2021 22:34:40 +0800 Subject: [PATCH 5/6] revert and keep one case --- .../debezium/DebeziumMongoDbSourceTester.java | 6 ++-- .../debezium/DebeziumMsSqlSourceTester.java | 6 ++-- .../debezium/DebeziumMySqlSourceTester.java | 7 ++-- .../DebeziumPostgreSqlSourceTester.java | 6 ++-- .../debezium/PulsarDebeziumSourcesTest.java | 36 +++++-------------- 5 files changed, 17 insertions(+), 44 deletions(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java index 24eb491348f8b..110ff11c00d07 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/DebeziumMongoDbSourceTester.java @@ -39,7 +39,7 @@ public class DebeziumMongoDbSourceTester extends SourceTester confirmedFlushLsn = new AtomicReference<>("not read yet"); - public DebeziumPostgreSqlSourceTester(PulsarCluster cluster, boolean testWithClientBuilder) { + public DebeziumPostgreSqlSourceTester(PulsarCluster cluster) { super(NAME); this.pulsarCluster = cluster; /* @@ -80,9 +80,7 @@ public DebeziumPostgreSqlSourceTester(PulsarCluster cluster, boolean testWithCli sourceConfig.put("database.dbname", "postgres"); sourceConfig.put("schema.whitelist", "inventory"); sourceConfig.put("table.blacklist", "inventory.spatial_ref_sys,inventory.geom"); - if (!testWithClientBuilder) { - sourceConfig.put("database.history.pulsar.service.url", pulsarServiceUrl); - } + sourceConfig.put("database.history.pulsar.service.url", pulsarServiceUrl); sourceConfig.put("topic.namespace", "debezium/postgresql"); } diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java index e5236ee16e522..246b52b5dd9fc 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sources/debezium/PulsarDebeziumSourcesTest.java @@ -61,32 +61,18 @@ public void testDebeziumMySqlSourceAvro() throws Exception { @Test(groups = "source") public void testDebeziumPostgreSqlSource() throws Exception { - testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true); } - @Test(groups = "source") - public void testDebeziumPostgreSqlSourceWithClientBuilder() throws Exception { - testDebeziumPostgreSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, true); - } @Test(groups = "source") public void testDebeziumMongoDbSource() throws Exception{ - testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true, false); - } - - @Test(groups = "source") - public void testDebeziumMongoDbSourceWithClientBuilder() throws Exception{ - testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true, true); + testDebeziumMongoDbConnect("org.apache.kafka.connect.json.JsonConverter", true); } @Test(groups = "source") public void testDebeziumMsSqlSource() throws Exception{ - testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); - } - - @Test(groups = "source") - public void testDebeziumMsSqlSourceWithClientBuilder() throws Exception{ - testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true, false); + testDebeziumMsSqlConnect("org.apache.kafka.connect.json.JsonConverter", true); } private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWithEnvelope, @@ -139,8 +125,7 @@ private void testDebeziumMySqlConnect(String converterClassName, boolean jsonWit runner.testSource(sourceTester); } - private void testDebeziumPostgreSqlConnect(String converterClassName, boolean jsonWithEnvelope, - boolean testWithClientBuilder) throws Exception { + private void testDebeziumPostgreSqlConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -164,7 +149,7 @@ private void testDebeziumPostgreSqlConnect(String converterClassName, boolean js admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumPostgreSqlSourceTester sourceTester = new DebeziumPostgreSqlSourceTester(pulsarCluster, testWithClientBuilder); + DebeziumPostgreSqlSourceTester sourceTester = new DebeziumPostgreSqlSourceTester(pulsarCluster); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); // setup debezium postgresql server @@ -178,8 +163,7 @@ private void testDebeziumPostgreSqlConnect(String converterClassName, boolean js runner.testSource(sourceTester); } - private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonWithEnvelope, - boolean testWithClientBuilder) throws Exception { + private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -204,8 +188,7 @@ private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonW admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumMongoDbSourceTester sourceTester = - new DebeziumMongoDbSourceTester(pulsarCluster, testWithClientBuilder); + DebeziumMongoDbSourceTester sourceTester = new DebeziumMongoDbSourceTester(pulsarCluster); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); // setup debezium mongodb server @@ -219,8 +202,7 @@ private void testDebeziumMongoDbConnect(String converterClassName, boolean jsonW runner.testSource(sourceTester); } - private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWithEnvelope, - boolean testWithClientBuilder) throws Exception { + private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWithEnvelope) throws Exception { final String tenant = TopicName.PUBLIC_TENANT; final String namespace = TopicName.DEFAULT_NAMESPACE; @@ -243,7 +225,7 @@ private void testDebeziumMsSqlConnect(String converterClassName, boolean jsonWit admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup - DebeziumMsSqlSourceTester sourceTester = new DebeziumMsSqlSourceTester(pulsarCluster, testWithClientBuilder); + DebeziumMsSqlSourceTester sourceTester = new DebeziumMsSqlSourceTester(pulsarCluster); sourceTester.getSourceConfig().put("json-with-envelope", jsonWithEnvelope); DebeziumMsSqlContainer msSqlContainer = new DebeziumMsSqlContainer(pulsarCluster.getClusterName()); From 47a87f6d6820bb8c6167d5333005f314595726ef Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Fri, 19 Nov 2021 09:22:15 +0800 Subject: [PATCH 6/6] fix CI --- .../io/sources/debezium/DebeziumMySqlSourceTester.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) 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 3cb64db8a7de0..7958fa019925f 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