From fb5acf51f78cc43e5c0496c90ec95f3e328c5ecb Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 19 Nov 2021 09:28:59 -0500 Subject: [PATCH] GH-182: Upgrade Debezium to 1.7.1 Fixes https://github.com/spring-cloud/stream-applications/issues/182 * Upgrade Debezium dependency to 1.7.1 * Remove explicit dependencies for DBs and remove excludes for them from Debezium deps; rely fully on whatever Debezium connectors bring for us * Fix `EmbeddedEngine` for compatibility with the Debezium 1.7.1 * Fix CDC tests for the current state of Debezium results * Fix Checkstyle violations in the MQTT modules --- applications/sink/mqtt-sink/README.adoc | 1 + .../source/cdc-debezium-source/pom.xml | 40 ---------------- .../cdc/CdcFlatteningIntegrationTest.java | 15 ++++-- .../app/source/cdc/CdcMySqlTestSupport.java | 4 +- .../CdcSourceDatabasesIntegrationTest.java | 10 ++-- .../mysql_insert_inventory_products_106.json | 2 +- applications/source/mqtt-source/README.adoc | 1 + .../common/cdc-debezium-boot-starter/pom.xml | 45 ++++++----------- .../cdc/CdcBootStarterIntegrationTest.java | 2 +- functions/common/cdc-debezium-common/pom.xml | 42 +--------------- .../cloud/fn/common/cdc/EmbeddedEngine.java | 42 +++++++++++++++- .../fn/consumer/mqtt/MqttConsumerTests.java | 1 - .../supplier/cdc-debezium-supplier/pom.xml | 48 ++----------------- .../fn/supplier/mqtt/MqttSupplierTests.java | 7 ++- 14 files changed, 84 insertions(+), 176 deletions(-) diff --git a/applications/sink/mqtt-sink/README.adoc b/applications/sink/mqtt-sink/README.adoc index 2648234a..61ccba5b 100644 --- a/applications/sink/mqtt-sink/README.adoc +++ b/applications/sink/mqtt-sink/README.adoc @@ -24,6 +24,7 @@ $$keep-alive-interval$$:: $$the ping interval in seconds.$$ *($$Integer$$, defau $$password$$:: $$the password to use when connecting to the broker.$$ *($$String$$, default: `$$guest$$`)* $$persistence$$:: $$'memory' or 'file'.$$ *($$String$$, default: `$$memory$$`)* $$persistence-directory$$:: $$Persistence directory.$$ *($$String$$, default: `$$/tmp/paho$$`)* +$$ssl-properties$$:: $$MQTT Client SSL properties.$$ *($$Map$$, default: `$$$$`)* $$url$$:: $$location of the mqtt broker(s) (comma-delimited list).$$ *($$String[]$$, default: `$$[tcp://localhost:1883]$$`)* $$username$$:: $$the username to use when connecting to the broker.$$ *($$String$$, default: `$$guest$$`)* diff --git a/applications/source/cdc-debezium-source/pom.xml b/applications/source/cdc-debezium-source/pom.xml index 31a924bd..99536786 100644 --- a/applications/source/cdc-debezium-source/pom.xml +++ b/applications/source/cdc-debezium-source/pom.xml @@ -26,21 +26,6 @@ org.springframework.cloud.fn cdc-debezium-supplier ${java-functions.version} - - - slf4j-log4j12 - org.slf4j - - - mysql - mysql-connector-java - - - - - mysql - mysql-connector-java - ${mysql-connector-java.version} org.springframework.cloud.fn @@ -60,10 +45,6 @@ ${stream-apps-core.version} test - - org.springframework.kafka - spring-kafka - org.springframework spring-jdbc @@ -105,27 +86,6 @@ test - - org.mongodb - mongodb-driver - 3.12.7 - test - - - - org.mongodb - bson - 3.12.7 - test - - - - org.mongodb - mongodb-driver-core - 3.12.7 - test - - diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java index cfbd4180..7d863a63 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2020 the original author or authors. + * Copyright 2020-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.app.source.cdc; import java.util.List; +import net.javacrumbs.jsonunit.core.Configuration; import org.junit.jupiter.api.Test; import org.springframework.boot.test.context.FilteredClassLoader; @@ -42,6 +43,7 @@ import static org.springframework.cloud.stream.app.source.cdc.CdcTestUtils.resou /** * @author Christian Tzolov * @author David Turanski + * @author Artem Bilan */ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { @@ -85,12 +87,14 @@ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { assertJsonEquals(resourceToString( "classpath:/json/mysql_ddl_drop_inventory_address_table.json"), - toString(messages.get(1).getPayload())); + toString(messages.get(1).getPayload()), + Configuration.empty().whenIgnoringPaths("schemaName", "tableChanges", "source.sequence", "source.ts_ms")); assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector"); assertJsonEquals("{\"databaseName\":\"inventory\"}", toString(messages.get(1).getHeaders().get("cdc_key"))); assertJsonEquals(resourceToString("classpath:/json/mysql_insert_inventory_products_106.json"), - toString(messages.get(39).getPayload())); + toString(messages.get(39).getPayload()), + Configuration.empty().whenIgnoringPaths("source.sequence", "source.ts_ms")); assertThat(messages.get(39).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.products"); assertJsonEquals("{\"id\":106}", toString(messages.get(39).getHeaders().get("cdc_key"))); @@ -172,7 +176,8 @@ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { assertJsonEquals(resourceToString( "classpath:/json/mysql_ddl_drop_inventory_address_table.json"), - toString(messages.get(1).getPayload())); + toString(messages.get(1).getPayload()), + Configuration.empty().whenIgnoringPaths("schemaName", "tableChanges", "source.sequence", "source.ts_ms")); assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector"); assertJsonEquals("{\"databaseName\":\"inventory\"}", toString(messages.get(1).getHeaders().get("cdc_key"))); @@ -189,7 +194,7 @@ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { assertJsonEquals("{\"id\":106}", toString(messages.get(39).getHeaders().get("cdc_key"))); if (flatteningProps.isEnabled() && flatteningProps.getAddHeaders().contains("op")) { - assertThat(messages.get(39).getHeaders().get("__op")).isEqualTo("c"); + assertThat(messages.get(39).getHeaders().get("__op")).isEqualTo("r"); } jdbcTemplate.update( diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java index 55425bde..00c68f15 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2020 the original author or authors. + * Copyright 2020-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -33,7 +33,7 @@ public abstract class CdcMySqlTestSupport { static String MAPPED_PORT; - static GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.3") + static GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.7.1.Final") .withEnv("MYSQL_ROOT_PASSWORD", "debezium") .withEnv("MYSQL_USER", "mysqluser") .withEnv("MYSQL_PASSWORD", "mysqlpw") diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java index dffa36a7..add32b25 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java @@ -52,7 +52,7 @@ public class CdcSourceDatabasesIntegrationTest { TestChannelBinderConfiguration.getCompleteConfiguration(TestCdcSourceApplication.class)) .web(WebApplicationType.NONE) .properties("spring.cloud.stream.function.definition=cdcSupplier", - "cdc.name=my-sql-connector", + "cdc.name=my-connector", "cdc.flattening.dropTombstones=false", "cdc.schema=false", "cdc.flattening.enabled=true", @@ -63,7 +63,7 @@ public class CdcSourceDatabasesIntegrationTest { @Test public void mysql() { - GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.3") + GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.7.1.Final") .withEnv("MYSQL_ROOT_PASSWORD", "debezium") .withEnv("MYSQL_USER", "mysqluser") .withEnv("MYSQL_PASSWORD", "mysqlpw") @@ -124,7 +124,7 @@ public class CdcSourceDatabasesIntegrationTest { @Test public void postgres() { - GenericContainer postgres = new GenericContainer("debezium/example-postgres:1.3") + GenericContainer postgres = new GenericContainer("debezium/example-postgres:1.7.1.Final") .withEnv("POSTGRES_USER", "postgres") .withEnv("POSTGRES_PASSWORD", "postgres") .withExposedPorts(5432); @@ -157,7 +157,7 @@ public class CdcSourceDatabasesIntegrationTest { @Test @Disabled public void mongodb() { - GenericContainer mongodb = new GenericContainer("debezium/example-mongodb:1.3") + GenericContainer mongodb = new GenericContainer("debezium/example-mongodb:1.7.1.Final") .withEnv("MONGODB_USER", "debezium") .withEnv("MONGODB_PASSWORD", "dbz") .withExposedPorts(27017); @@ -169,7 +169,7 @@ public class CdcSourceDatabasesIntegrationTest { "--cdc.config.mongodb.name=dbserver1", "--cdc.config.mongodb.user=debezium", "--cdc.config.mongodb.password=dbz", - "--cdc.config.database.whitelist=inventory")) { + "--cdc.config.collection.include.list=inventory[.]*")) { OutputDestination outputDestination = context.getBean(OutputDestination.class); // Using local region here List> messages = receiveAll(outputDestination); diff --git a/applications/source/cdc-debezium-source/src/test/resources/json/mysql_insert_inventory_products_106.json b/applications/source/cdc-debezium-source/src/test/resources/json/mysql_insert_inventory_products_106.json index 363f33c3..3806280c 100644 --- a/applications/source/cdc-debezium-source/src/test/resources/json/mysql_insert_inventory_products_106.json +++ b/applications/source/cdc-debezium-source/src/test/resources/json/mysql_insert_inventory_products_106.json @@ -22,7 +22,7 @@ "thread": null, "query": null }, - "op": "c", + "op": "r", "ts_ms": "${json-unit.ignore}", "transaction": null } diff --git a/applications/source/mqtt-source/README.adoc b/applications/source/mqtt-source/README.adoc index e677727a..ed098c9b 100644 --- a/applications/source/mqtt-source/README.adoc +++ b/applications/source/mqtt-source/README.adoc @@ -24,6 +24,7 @@ $$keep-alive-interval$$:: $$the ping interval in seconds.$$ *($$Integer$$, defau $$password$$:: $$the password to use when connecting to the broker.$$ *($$String$$, default: `$$guest$$`)* $$persistence$$:: $$'memory' or 'file'.$$ *($$String$$, default: `$$memory$$`)* $$persistence-directory$$:: $$Persistence directory.$$ *($$String$$, default: `$$/tmp/paho$$`)* +$$ssl-properties$$:: $$MQTT Client SSL properties.$$ *($$Map$$, default: `$$$$`)* $$url$$:: $$location of the mqtt broker(s) (comma-delimited list).$$ *($$String[]$$, default: `$$[tcp://localhost:1883]$$`)* $$username$$:: $$the username to use when connecting to the broker.$$ *($$String$$, default: `$$guest$$`)* diff --git a/functions/common/cdc-debezium-boot-starter/pom.xml b/functions/common/cdc-debezium-boot-starter/pom.xml index 80c3cb71..2b618c90 100644 --- a/functions/common/cdc-debezium-boot-starter/pom.xml +++ b/functions/common/cdc-debezium-boot-starter/pom.xml @@ -15,7 +15,7 @@ Change Data Capture (CDC) Debezium Boot Starter - 1.3.0.Final + 1.7.1.Final @@ -38,24 +38,24 @@ debezium-connector-mysql - mysql-connector-java - mysql + slf4j-log4j12 + org.slf4j ${version.debezium} - - - - - - - - - - - + + io.debezium + debezium-connector-mongodb + + + slf4j-log4j12 + org.slf4j + + + ${version.debezium} + io.debezium @@ -91,23 +91,6 @@ ${version.debezium} - - mysql - mysql-connector-java - 8.0.13 - - - - org.testcontainers - junit-jupiter - test - - - - org.testcontainers - mysql - - org.springframework spring-jdbc diff --git a/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java b/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java index 63951811..58ef862a 100644 --- a/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java +++ b/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java @@ -50,7 +50,7 @@ public class CdcBootStarterIntegrationTest { @Container static GenericContainer debeziumMySQL = - new GenericContainer<>(DockerImageName.parse("debezium/example-mysql:1.3")) + new GenericContainer<>(DockerImageName.parse("debezium/example-mysql:1.7.1.Final")) .withEnv("MYSQL_ROOT_PASSWORD", "debezium") .withEnv("MYSQL_USER", "mysqluser") .withEnv("MYSQL_PASSWORD", "mysqlpw") diff --git a/functions/common/cdc-debezium-common/pom.xml b/functions/common/cdc-debezium-common/pom.xml index d3bef15e..0a1092b8 100644 --- a/functions/common/cdc-debezium-common/pom.xml +++ b/functions/common/cdc-debezium-common/pom.xml @@ -15,12 +15,7 @@ Change Data Capture (CDC) Debezium Common - 1.3.1.Final - - 42.2.5 - 5.2.1.RELEASE - - 5.1.46 + 1.7.1.Final @@ -39,29 +34,17 @@ io.debezium debezium-connector-mysql - - - mysql-connector-java - mysql - - true ${version.debezium} - - mysql - mysql-connector-java - true - ${mysql-connector-java} - - io.debezium debezium-connector-postgres true ${version.debezium} + io.debezium debezium-connector-mongodb @@ -87,16 +70,6 @@ ${project.version} - - - org.springframework.boot - spring-boot-starter-json - - - org.springframework.boot - spring-boot-starter-validation - - org.springframework.integration spring-integration-ip @@ -107,17 +80,6 @@ ${revision} - - org.springframework.boot - spring-boot-configuration-processor - provided - - - org.springframework.boot - spring-boot-starter-test - test - - diff --git a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/EmbeddedEngine.java b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/EmbeddedEngine.java index 624653be..687275c0 100644 --- a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/EmbeddedEngine.java +++ b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/EmbeddedEngine.java @@ -8,6 +8,7 @@ package org.springframework.cloud.fn.common.cdc; import java.io.IOException; import java.time.Duration; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -460,6 +461,31 @@ public final class EmbeddedEngine implements DebeziumEngine { public static interface ChangeConsumer extends DebeziumEngine.ChangeConsumer { } + protected class SourceRecordOffsets implements DebeziumEngine.Offsets { + + private final HashMap offsets = new HashMap<>(); + + /** + * Performs {@link HashMap#put(Object, Object)} on the offsets map. + * + * @param key key with which to put the value + * @param value value to be put with the key + */ + @Override + public void set(String key, Object value) { + offsets.put(key, value); + } + + /** + * Retrieves the offsets map. + * + * @return HashMap of the offsets + */ + protected HashMap getOffsets() { + return offsets; + } + } + private static ChangeConsumer buildDefaultChangeConsumer(Consumer consumer) { return new ChangeConsumer() { @@ -912,9 +938,23 @@ public final class EmbeddedEngine implements DebeziumEngine { } @Override - public synchronized void markBatchFinished() { + public synchronized void markBatchFinished() throws InterruptedException { maybeFlush(offsetWriter, offsetCommitPolicy, commitTimeout, task); } + + @Override + public synchronized void markProcessed(SourceRecord record, DebeziumEngine.Offsets sourceOffsets) throws InterruptedException { + SourceRecordOffsets offsets = (SourceRecordOffsets) sourceOffsets; + SourceRecord recordWithUpdatedOffsets = new SourceRecord(record.sourcePartition(), offsets.getOffsets(), record.topic(), + record.kafkaPartition(), record.keySchema(), record.key(), record.valueSchema(), record.value(), + record.timestamp(), record.headers()); + markProcessed(recordWithUpdatedOffsets); + } + + @Override + public DebeziumEngine.Offsets buildOffsets() { + return new SourceRecordOffsets(); + } }; } diff --git a/functions/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java b/functions/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java index dad13103..752b5dfa 100644 --- a/functions/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java +++ b/functions/consumer/mqtt-consumer/src/test/java/org/springframework/cloud/fn/consumer/mqtt/MqttConsumerTests.java @@ -35,7 +35,6 @@ import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.mqtt.core.MqttPahoClientFactory; import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter; import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter; -import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.test.annotation.DirtiesContext; diff --git a/functions/supplier/cdc-debezium-supplier/pom.xml b/functions/supplier/cdc-debezium-supplier/pom.xml index 95a97976..737898fd 100644 --- a/functions/supplier/cdc-debezium-supplier/pom.xml +++ b/functions/supplier/cdc-debezium-supplier/pom.xml @@ -14,8 +14,7 @@ CDC Debezium Suppliers - 1.3.1.Final - 8.0.13 + 1.7.1.Final @@ -24,25 +23,17 @@ cdc-debezium-common ${project.version} - - mysql - mysql-connector-java - ${mysql.version} - - - io.debezium debezium-connector-mysql - mysql-connector-java - mysql + slf4j-log4j12 + org.slf4j ${version.debezium} - io.debezium debezium-connector-mongodb @@ -87,38 +78,5 @@ ${version.debezium} - - - - org.springframework.boot - spring-boot-starter-integration - - - - org.springframework.boot - spring-boot-starter-validation - - - org.springframework.boot - spring-boot-configuration-processor - provided - - - - org.springframework.boot - spring-boot-starter-test - test - - - io.projectreactor - reactor-test - test - - - org.springframework.integration - spring-integration-test - test - - diff --git a/functions/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java b/functions/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java index 21c66ece..47a3ec59 100644 --- a/functions/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java +++ b/functions/supplier/mqtt-supplier/src/test/java/org/springframework/cloud/fn/supplier/mqtt/MqttSupplierTests.java @@ -16,8 +16,6 @@ package org.springframework.cloud.fn.supplier.mqtt; -import static org.assertj.core.api.Assertions.assertThat; - import java.util.Properties; import java.util.function.Supplier; @@ -27,6 +25,8 @@ import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.testcontainers.containers.GenericContainer; +import reactor.core.publisher.Flux; +import reactor.test.StepVerifier; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.SpringBootApplication; @@ -40,8 +40,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.support.MessageBuilder; import org.springframework.test.annotation.DirtiesContext; -import reactor.core.publisher.Flux; -import reactor.test.StepVerifier; +import static org.assertj.core.api.Assertions.assertThat; /** * Tests for Mqtt Supplier.