From a9692efc75eb571275c2f188bece56e8d86c2f8a Mon Sep 17 00:00:00 2001 From: Christian Tzolov Date: Fri, 4 Dec 2020 19:42:10 +0100 Subject: [PATCH] Fix CDC issues - Fix cdc_key header string wrapping issue that causes problems with the spring.cloud.stream.kafka.default.producer.messageKeyExpression. - Rename Flattering property to the correct Flattening. - Remove wrong dependencies. --- .../source/cdc-debezium-source/README.adoc | 27 +++--- .../source/cdc-debezium-source/pom.xml | 11 +-- ...onfiguration-metadata-whitelist.properties | 2 +- ...dataflow-configuration-metadata.properties | 2 +- .../src/main/resources/application.properties | 2 - .../cdc/CdcDeleteHandlingIntegrationTest.java | 20 ++--- ...java => CdcFlatteningIntegrationTest.java} | 82 +++++++++---------- ...tSupport.java => CdcMySqlTestSupport.java} | 4 +- .../CdcSourceDatabasesIntegrationTest.java | 51 +++++++----- ...ttened_insert_inventory_products_106.json} | 0 ...flattened_update_inventory_customers.json} | 0 applications/stream-applications-core/pom.xml | 2 +- .../cdc-debezium-boot-starter/README.adoc | 14 ++-- .../fn/common/cdc/CdcAutoConfiguration.java | 4 +- .../cdc/CdcBootStarterIntegrationTest.java | 6 +- functions/common/cdc-debezium-common/pom.xml | 2 +- .../fn/common/cdc/CdcCommonConfiguration.java | 16 ++-- .../fn/common/cdc/CdcCommonProperties.java | 12 +-- .../cloud/fn/common/cdc/EmbeddedEngine.java | 16 +--- .../src/main/resources/application.properties | 2 +- .../supplier/cdc-debezium-supplier/pom.xml | 2 +- .../cdc/CdcSupplierConfiguration.java | 11 +-- .../src/main/resources/application.properties | 2 +- 23 files changed, 142 insertions(+), 148 deletions(-) delete mode 100644 applications/source/cdc-debezium-source/src/main/resources/application.properties rename applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/{CdcFlatteringIntegrationTest.java => CdcFlatteningIntegrationTest.java} (81%) rename applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/{CdcTestSupport.java => CdcMySqlTestSupport.java} (95%) rename applications/source/cdc-debezium-source/src/test/resources/json/{mysql_flattered_insert_inventory_products_106.json => mysql_flattened_insert_inventory_products_106.json} (100%) rename applications/source/cdc-debezium-source/src/test/resources/json/{mysql_flattered_update_inventory_customers.json => mysql_flattened_update_inventory_customers.json} (100%) diff --git a/applications/source/cdc-debezium-source/README.adoc b/applications/source/cdc-debezium-source/README.adoc index 46ce9008..4a232583 100644 --- a/applications/source/cdc-debezium-source/README.adoc +++ b/applications/source/cdc-debezium-source/README.adoc @@ -17,8 +17,7 @@ The CDC Source introduces a new default `BackingOffsetStore` configuration, base == Options //tag::configuration-properties[] -Properties grouped by prefix: - +Properties grouped by a prefix: === cdc @@ -27,13 +26,13 @@ $$connector$$:: $$Shortcut for the cdc.config.connector.class property. Either o $$name$$:: $$Unique name for this sourceConnector instance.$$ *($$String$$, default: `$$$$`)* $$schema$$:: $$Include the schema's as part of the outbound message.$$ *($$Boolean$$, default: `$$false$$`)* -=== cdc.flattering +=== cdc.flattening $$add-fields$$:: $$Comma separated list of metadata fields to add to the flattened message. The fields will be prefixed with "__" or "__[<]struct]__", depending on the specification of the struct.$$ *($$String$$, default: `$$$$`)* $$add-headers$$:: $$Comma separated list specify a list of metadata fields to add to the header of the flattened message. The fields will be prefixed with "__" or "__[struct]__".$$ *($$String$$, default: `$$$$`)* $$delete-handling-mode$$:: $$Options for handling deleted records: (1) none - pass the records through, (2) drop - remove the records and (3) rewrite - add a '__deleted' field to the records.$$ *($$DeleteHandlingMode$$, default: `$$$$`, possible values: `drop`,`rewrite`,`none`)* $$drop-tombstones$$:: $$By default Debezium generates tombstone records to enable Kafka compaction on deleted records. The dropTombstones can suppress the tombstone records.$$ *($$Boolean$$, default: `$$true$$`)* -$$enabled$$:: $$Enable flattering the source record events (https://debezium.io/docs/configuration/event-flattening).$$ *($$Boolean$$, default: `$$true$$`)* +$$enabled$$:: $$Enable flattening the source record events (https://debezium.io/docs/configuration/event-flattening).$$ *($$Boolean$$, default: `$$true$$`)* === cdc.offset @@ -117,11 +116,11 @@ The table below lists all available shortcuts along with the Debezium properties |cdc.config.offset.storage |`metadata` : MetadataStoreOffsetBackingStore, `file` : FileOffsetBackingStore, `kafka` : KafkaOffsetBackingStore, `memory` : MemoryOffsetBackingStore -|cdc.flattering.drop-tombstones +|cdc.flattening.drop-tombstones |cdc.config.drop.tombstones | -|cdc.flattering.delete-handling-mode +|cdc.flattening.delete-handling-mode |cdc.config.delete.handling.mode |`none` : none, `drop` : drop, `rewrite` : rewrite @@ -133,7 +132,7 @@ The `CDC Source` uses the Debezium utilities, and currently supports CDC for fiv == Examples and Testing -The [CdcSourceIntegrationTest](), [CdcDeleteHandlingIntegrationTest]() and [CdcFlatteringIntegrationTest]() integration tests use test databases fixtures, running on the local machine. +The [CdcSourceIntegrationTest](), [CdcDeleteHandlingIntegrationTest]() and [CdcFlatteningIntegrationTest]() integration tests use test databases fixtures, running on the local machine. We use pre-build debezium docker database images. The Maven builds create the test databases fixtures with the help of the `docker-maven-plugin`. @@ -174,13 +173,13 @@ cdc.config.database.hostname=localhost # <3> cdc.config.database.port=3306 # <3> cdc.schema=true # <4> -cdc.flattering.enabled=true # <5> +cdc.flattening.enabled=true # <5> ---- <1> Configures the CDC Source to use https://debezium.io/docs/connectors/mysql/[MySqlConnector]. (equivalent to setting `cdc.config.connector.class=io.debezium.connector.mysql.MySqlConnector`). <2> Metadata used to identify and dispatch the incoming events. <3> Connection to the MySQL server running on `localhost:3306` as `debezium` user. <4> Includes the https://debezium.io/docs/connectors/mysql/#change-events-value[Change Event Value] schema in the `SourceRecord` events. -<5> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattering]. +<5> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattening]. You can run also the `CdcSourceIntegrationTests#CdcMysqlTests` using this mysql configuration. @@ -216,7 +215,7 @@ cdc.config.database.hostname=localhost # <4> cdc.config.database.port=5432 # <4> cdc.schema=true # <5> -cdc.flattering.enabled=true # <6> +cdc.flattening.enabled=true # <6> ---- <1> Configures `CDC Source` to use https://debezium.io/docs/connectors/postgresql/[PostgresConnector]. Equivalent for setting `cdc.config.connector.class=io.debezium.connector.postgresql.PostgresConnector`. @@ -224,7 +223,7 @@ cdc.flattering.enabled=true # <6> <3> Metadata used to identify and dispatch the incoming events. <4> Connection to the PostgreSQL server running on `localhost:5432` as `postgres` user. <5> Includes the https://debezium.io/docs/connectors/mysql/#change-events-value[Change Event Value] schema in the `SourceRecord` events. -<6> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattering]. +<6> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattening]. You can run also the `CdcSourceIntegrationTests#CdcPostgresTests` using this mysql configuration. @@ -267,7 +266,7 @@ cdc.config.database.whitelist=inventory # <3> cdc.config.tasks.max=1 # <4> cdc.schema=true # <5> -cdc.flattering.enabled=true # <6> +cdc.flattening.enabled=true # <6> ---- <1> Configures `CDC Source` to use https://debezium.io/docs/connectors/mongodb/[MongoDB Connector]. This maps into `cdc.config.connector.class=io.debezium.connector.mongodb.MongodbSourceConnector`. @@ -275,7 +274,7 @@ cdc.flattering.enabled=true # <6> <3> Connection to the MongoDB running on `localhost:27017` as `debezium` user. <4> https://debezium.io/docs/connectors/mongodb/#tasks <5> Includes the https://debezium.io/docs/connectors/mysql/#change-events-value[Change Event Value] schema in the `SourceRecord` events. -<6> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattering]. +<6> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattening]. You can run also the `CdcSourceIntegrationTests#CdcPostgresTests` using this mysql configuration. @@ -336,6 +335,6 @@ cat ./inventory.sql | docker exec -i dbz_oracle sqlplus debezium/dbz@//localhost == Run standalone ``` -java -jar cdc-debezium-source.jar --cdc.connector=mysql --cdc.name=my-sql-connector --cdc.config.database.server.id=85744 --cdc.config.database.server.name=my-app-connector --cdc.config.database.user=debezium --cdc.config.database.password=dbz --cdc.config.database.hostname=localhost --cdc.config.database.port=3306 --cdc.schema=true --cdc.flattering.enabled=true +java -jar cdc-debezium-source.jar --cdc.connector=mysql --cdc.name=my-sql-connector --cdc.config.database.server.id=85744 --cdc.config.database.server.name=my-app-connector --cdc.config.database.user=debezium --cdc.config.database.password=dbz --cdc.config.database.hostname=localhost --cdc.config.database.port=3306 --cdc.schema=true --cdc.flattening.enabled=true ``` diff --git a/applications/source/cdc-debezium-source/pom.xml b/applications/source/cdc-debezium-source/pom.xml index b0cd87c6..bd3bdbf3 100644 --- a/applications/source/cdc-debezium-source/pom.xml +++ b/applications/source/cdc-debezium-source/pom.xml @@ -94,13 +94,6 @@ mysql - - org.springframework.cloud.fn - cdc-debezium-boot-starter - 1.0.0-M4 - test - - @@ -121,6 +114,10 @@ cdcSupplier + + org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration + headers['cdc_key'] + diff --git a/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties b/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties index 30b9254a..fa21a45a 100644 --- a/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties +++ b/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata-whitelist.properties @@ -1,7 +1,7 @@ configuration-properties.classes=org.springframework.cloud.fn.supplier.cdc.CdcSupplierProperties, \ org.springframework.cloud.fn.supplier.cdc.CdcSupplierProperties$Header, \ org.springframework.cloud.fn.common.cdc.CdcCommonProperties, \ - org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Flattering, \ + org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Flattening, \ org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Offset, \ org.springframework.cloud.fn.common.metadata.store.MetadataStoreProperties, \ org.springframework.cloud.fn.common.metadata.store.MetadataStoreProperties$Gemfire, \ diff --git a/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata.properties b/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata.properties index 30b9254a..fa21a45a 100644 --- a/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata.properties +++ b/applications/source/cdc-debezium-source/src/main/resources/META-INF/dataflow-configuration-metadata.properties @@ -1,7 +1,7 @@ configuration-properties.classes=org.springframework.cloud.fn.supplier.cdc.CdcSupplierProperties, \ org.springframework.cloud.fn.supplier.cdc.CdcSupplierProperties$Header, \ org.springframework.cloud.fn.common.cdc.CdcCommonProperties, \ - org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Flattering, \ + org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Flattening, \ org.springframework.cloud.fn.common.cdc.CdcCommonProperties$Offset, \ org.springframework.cloud.fn.common.metadata.store.MetadataStoreProperties, \ org.springframework.cloud.fn.common.metadata.store.MetadataStoreProperties$Gemfire, \ diff --git a/applications/source/cdc-debezium-source/src/main/resources/application.properties b/applications/source/cdc-debezium-source/src/main/resources/application.properties deleted file mode 100644 index 8fb76f9d..00000000 --- a/applications/source/cdc-debezium-source/src/main/resources/application.properties +++ /dev/null @@ -1,2 +0,0 @@ -spring.autoconfigure.exclude=org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration -spring.cloud.stream.kafka.default.producer.messageKeyExpression=headers['cdc_key'] diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcDeleteHandlingIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcDeleteHandlingIntegrationTest.java index a1824759..57b6ae5b 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcDeleteHandlingIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcDeleteHandlingIntegrationTest.java @@ -45,7 +45,7 @@ import static org.springframework.cloud.stream.app.source.cdc.CdcTestUtils.recei */ @Testcontainers -public class CdcDeleteHandlingIntegrationTest extends CdcTestSupport { +public class CdcDeleteHandlingIntegrationTest extends CdcMySqlTestSupport { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() .withUserConfiguration( @@ -54,7 +54,7 @@ public class CdcDeleteHandlingIntegrationTest extends CdcTestSupport { "spring.cloud.function.definition=cdcSupplier", "cdc.name=my-sql-connector", "cdc.schema=false", - "cdc.flattering.enabled=true", + "cdc.flattening.enabled=true", "cdc.stream.header.offset=true", "cdc.connector=mysql", "cdc.config.database.user=debezium", @@ -67,12 +67,12 @@ public class CdcDeleteHandlingIntegrationTest extends CdcTestSupport { @ParameterizedTest @ValueSource(strings = { - "cdc.flattering.deleteHandlingMode=none,cdc.flattering.dropTombstones=true", - "cdc.flattering.deleteHandlingMode=none,cdc.flattering.dropTombstones=false", - "cdc.flattering.deleteHandlingMode=drop,cdc.flattering.dropTombstones=true", - "cdc.flattering.deleteHandlingMode=drop,cdc.flattering.dropTombstones=false", - "cdc.flattering.deleteHandlingMode=rewrite,cdc.flattering.dropTombstones=true", - "cdc.flattering.deleteHandlingMode=rewrite,cdc.flattering.dropTombstones=false" + "cdc.flattening.deleteHandlingMode=none,cdc.flattening.dropTombstones=true", + "cdc.flattening.deleteHandlingMode=none,cdc.flattening.dropTombstones=false", + "cdc.flattening.deleteHandlingMode=drop,cdc.flattening.dropTombstones=true", + "cdc.flattening.deleteHandlingMode=drop,cdc.flattening.dropTombstones=false", + "cdc.flattening.deleteHandlingMode=rewrite,cdc.flattening.dropTombstones=true", + "cdc.flattening.deleteHandlingMode=rewrite,cdc.flattening.dropTombstones=false" }) public void handleRecordDeletions(String properties) { contextRunner.withPropertyValues(properties.split(",")) @@ -94,8 +94,8 @@ public class CdcDeleteHandlingIntegrationTest extends CdcTestSupport { boolean isKafkaPresent = ClassUtils.isPresent(ORG_SPRINGFRAMEWORK_KAFKA_SUPPORT_KAFKA_NULL, context.getClassLoader()); - CdcCommonProperties.DeleteHandlingMode deleteHandlingMode = props.getFlattering().getDeleteHandlingMode(); - boolean isDropTombstones = props.getFlattering().isDropTombstones(); + CdcCommonProperties.DeleteHandlingMode deleteHandlingMode = props.getFlattening().getDeleteHandlingMode(); + boolean isDropTombstones = props.getFlattening().isDropTombstones(); jdbcTemplate.update( "insert into `customers`(`first_name`,`last_name`,`email`) VALUES('Test666', 'Test666', 'Test666@spring.org')"); diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteringIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java similarity index 81% rename from applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteringIntegrationTest.java rename to applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java index 3d97ef97..4e9672ac 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteringIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java @@ -43,7 +43,7 @@ import static org.springframework.cloud.stream.app.source.cdc.CdcTestUtils.resou * @author Christian Tzolov * @author David Turanski */ -public class CdcFlatteringIntegrationTest extends CdcTestSupport { +public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() .withUserConfiguration( @@ -63,19 +63,19 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { "cdc.config.database.history=io.debezium.relational.history.MemoryDatabaseHistory"); @Test - public void noFlatteredResponseNoKafka() { - contextRunner.withPropertyValues("cdc.flattering.enabled=false") + public void noFlattenedResponseNoKafka() { + contextRunner.withPropertyValues("cdc.flattening.enabled=false") .withClassLoader(new FilteredClassLoader(KafkaNull.class)) // Remove Kafka from the classpath - .run(noFlatteringTest); + .run(noFlatteningTest); } @Test - public void noFlatteredResponseWithKafka() { - contextRunner.withPropertyValues("cdc.flattering.enabled=false") - .run(noFlatteringTest); + public void noFlattenedResponseWithKafka() { + contextRunner.withPropertyValues("cdc.flattening.enabled=false") + .run(noFlatteningTest); } - final ContextConsumer noFlatteringTest = context -> { + final ContextConsumer noFlatteningTest = context -> { OutputDestination outputDestination = context.getBean(OutputDestination.class); boolean isKafkaPresent = ClassUtils.isPresent(ORG_SPRINGFRAMEWORK_KAFKA_SUPPORT_KAFKA_NULL, context.getClassLoader()); @@ -127,40 +127,40 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { }; @Test - public void flatteredResponseNoKafka() { + public void flattenedResponseNoKafka() { contextRunner.withPropertyValues( - "cdc.flattering.enabled=true", - "cdc.flattering.deleteHandlingMode=none", - "cdc.flattering.dropTombstones=false", - "cdc.flattering.addHeaders=op", - "cdc.flattering.addFields=name,db") + "cdc.flattening.enabled=true", + "cdc.flattening.deleteHandlingMode=none", + "cdc.flattening.dropTombstones=false", + "cdc.flattening.addHeaders=op", + "cdc.flattening.addFields=name,db") .withClassLoader(new FilteredClassLoader(KafkaNull.class)) // Remove Kafka from the classpath - .run(flatteringTest); + .run(flatteningTest); } @Test - public void flatteredResponseWithKafka() { + public void flattenedResponseWithKafka() { contextRunner.withPropertyValues( - "cdc.flattering.enabled=true", - "cdc.flattering.deleteHandlingMode=none", - "cdc.flattering.dropTombstones=false", - "cdc.flattering.addHeaders=op", - "cdc.flattering.addFields=name,db") - .run(flatteringTest); + "cdc.flattening.enabled=true", + "cdc.flattening.deleteHandlingMode=none", + "cdc.flattening.dropTombstones=false", + "cdc.flattening.addHeaders=op", + "cdc.flattening.addFields=name,db") + .run(flatteningTest); } @Test - public void flatteredResponseWithKafkaDropTombstone() { + public void flattenedResponseWithKafkaDropTombstone() { contextRunner.withPropertyValues( - "cdc.flattering.enabled=true", - "cdc.flattering.deleteHandlingMode=none", - "cdc.flattering.dropTombstones=true", - "cdc.flattering.addHeaders=op", - "cdc.flattering.addFields=name,db") - .run(flatteringTest); + "cdc.flattening.enabled=true", + "cdc.flattening.deleteHandlingMode=none", + "cdc.flattening.dropTombstones=true", + "cdc.flattening.addHeaders=op", + "cdc.flattening.addFields=name,db") + .run(flatteningTest); } - final ContextConsumer flatteringTest = context -> { + final ContextConsumer flatteningTest = context -> { OutputDestination outputDestination = context.getBean(OutputDestination.class); boolean isKafkaPresent = ClassUtils.isPresent(ORG_SPRINGFRAMEWORK_KAFKA_SUPPORT_KAFKA_NULL, context.getClassLoader()); @@ -168,7 +168,7 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { List> messages = receiveAll(outputDestination); assertThat(messages).hasSizeGreaterThanOrEqualTo(52); - CdcCommonProperties.Flattering flatteringProps = context.getBean(CdcCommonProperties.class).getFlattering(); + CdcCommonProperties.Flattening flatteningProps = context.getBean(CdcCommonProperties.class).getFlattening(); assertJsonEquals(resourceToString( "classpath:/json/mysql_ddl_drop_inventory_address_table.json"), @@ -177,8 +177,8 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { assertJsonEquals("{\"databaseName\":\"inventory\"}", messages.get(1).getHeaders().get("cdc_key")); - if (flatteringProps.isEnabled()) { - assertJsonEquals(resourceToString("classpath:/json/mysql_flattered_insert_inventory_products_106.json"), + if (flatteningProps.isEnabled()) { + assertJsonEquals(resourceToString("classpath:/json/mysql_flattened_insert_inventory_products_106.json"), toString(messages.get(39).getPayload())); } else { @@ -188,7 +188,7 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { assertThat(messages.get(39).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.products"); assertJsonEquals("{\"id\":106}", messages.get(39).getHeaders().get("cdc_key")); - if (flatteringProps.isEnabled() && flatteringProps.getAddHeaders().contains("op")) { + if (flatteningProps.isEnabled() && flatteningProps.getAddHeaders().contains("op")) { assertThat(messages.get(39).getHeaders().get("__op")).isEqualTo("c"); } @@ -201,28 +201,28 @@ public class CdcFlatteringIntegrationTest extends CdcTestSupport { messages = receiveAll(outputDestination); - assertThat(messages).hasSize((!flatteringProps.isDropTombstones() && isKafkaPresent) ? 4 : 3); + assertThat(messages).hasSize((!flatteningProps.isDropTombstones() && isKafkaPresent) ? 4 : 3); - assertJsonEquals(resourceToString("classpath:/json/mysql_flattered_update_inventory_customers.json"), + assertJsonEquals(resourceToString("classpath:/json/mysql_flattened_update_inventory_customers.json"), toString(messages.get(1).getPayload())); assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.customers"); assertJsonEquals("{\"id\":" + newRecordId + "}", messages.get(1).getHeaders().get("cdc_key")); - if (!StringUtils.isEmpty(flatteringProps.getAddHeaders()) && flatteringProps.getAddHeaders().contains("op")) { + if (!StringUtils.isEmpty(flatteningProps.getAddHeaders()) && flatteningProps.getAddHeaders().contains("op")) { assertThat(messages.get(1).getHeaders().get("__op")).isEqualTo("u"); } - if (flatteringProps.getDeleteHandlingMode() == CdcCommonProperties.DeleteHandlingMode.none) { + if (flatteningProps.getDeleteHandlingMode() == CdcCommonProperties.DeleteHandlingMode.none) { assertThat(toString(messages.get(2).getPayload())).isEqualTo("null"); assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.customers"); assertJsonEquals("{\"id\":" + newRecordId + "}", messages.get(1).getHeaders().get("cdc_key")); - if (!StringUtils.isEmpty(flatteringProps.getAddHeaders()) - && flatteringProps.getAddHeaders().contains("op")) { + if (!StringUtils.isEmpty(flatteningProps.getAddHeaders()) + && flatteningProps.getAddHeaders().contains("op")) { assertThat(messages.get(2).getHeaders().get("__op")).isEqualTo("d"); } } - if (!flatteringProps.isDropTombstones() && isKafkaPresent) { + if (!flatteningProps.isDropTombstones() && isKafkaPresent) { assertThat(messages.get(3).getPayload().getClass().getCanonicalName()) .isEqualTo(ORG_SPRINGFRAMEWORK_KAFKA_SUPPORT_KAFKA_NULL, "Tombstones event should have KafkaNull payload"); diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcTestSupport.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java similarity index 95% rename from applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcTestSupport.java rename to applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java index c879f2e0..03c9a72c 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcTestSupport.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java @@ -24,13 +24,13 @@ import org.springframework.jdbc.core.JdbcTemplate; /** * @author David Turanski */ -public abstract class CdcTestSupport { +public abstract class CdcMySqlTestSupport { static final String DATABASE_NAME = "inventory"; static String MAPPED_PORT; - static GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.0") + static GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.3") .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 b26eac51..5981fa51 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 @@ -23,7 +23,6 @@ import org.junit.jupiter.api.Test; import org.slf4j.LoggerFactory; import org.testcontainers.containers.GenericContainer; import org.testcontainers.containers.output.Slf4jLogConsumer; -import org.testcontainers.containers.wait.strategy.Wait; import org.testcontainers.images.builder.ImageFromDockerfile; import org.springframework.boot.WebApplicationType; @@ -41,23 +40,33 @@ import static org.springframework.cloud.stream.app.source.cdc.CdcTestUtils.recei * @author David Turanski */ @Disabled("Run as needed if there is an issue with a specific connector") -public class CdcSourceDatabasesIntegrationTest extends CdcTestSupport { +public class CdcSourceDatabasesIntegrationTest { private final SpringApplicationBuilder applicationBuilder = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(TestCdcSourceApplication.class)) - .web(WebApplicationType.NONE) - .properties("spring.cloud.stream.function.definition=cdcSupplier", - "cdc.name=my-sql-connector", - "cdc.flattering.dropTombstones=false", - "cdc.schema=false", - "cdc.flattering.enabled=true", - "cdc.stream.header.offset=true", - // "cdc.config.database.server.id=85744", - "cdc.config.database.server.name=my-app-connector", - "cdc.config.database.history=io.debezium.relational.history.MemoryDatabaseHistory"); + .web(WebApplicationType.NONE) + .properties("spring.cloud.stream.function.definition=cdcSupplier", + "cdc.name=my-sql-connector", + "cdc.flattening.dropTombstones=false", + "cdc.schema=false", + "cdc.flattening.enabled=true", + "cdc.stream.header.offset=true", + // "cdc.config.database.server.id=85744", + "cdc.config.database.server.name=my-app-connector", + "cdc.config.database.history=io.debezium.relational.history.MemoryDatabaseHistory"); @Test public void mysql() { + GenericContainer debeziumMySQL = new GenericContainer<>("debezium/example-mysql:1.3") + .withEnv("MYSQL_ROOT_PASSWORD", "debezium") + .withEnv("MYSQL_USER", "mysqluser") + .withEnv("MYSQL_PASSWORD", "mysqlpw") + // .withLogConsumer(new Slf4jLogConsumer(LoggerFactory.getLogger("mysql"))) + .withExposedPorts(3306); + debeziumMySQL.start(); + + String MAPPED_PORT = String.valueOf(debeziumMySQL.getMappedPort(3306)); + try (ConfigurableApplicationContext context = applicationBuilder .run("--cdc.connector=mysql", "--cdc.config.database.user=debezium", @@ -79,13 +88,15 @@ public class CdcSourceDatabasesIntegrationTest extends CdcTestSupport { .withFileFromClasspath("import-data.sh", "sqlserver/import-data.sh") .withFileFromClasspath("inventory.sql", "sqlserver/inventory.sql") .withFileFromClasspath("entrypoint.sh", "sqlserver/entrypoint.sh")) - .withEnv("ACCEPT_EULA", "Y") - .withEnv("MSSQL_PID", "Standard") - .withEnv("SA_PASSWORD", "Password!") - .withEnv("MSSQL_AGENT_ENABLED", "true") - .withLogConsumer(new Slf4jLogConsumer(LoggerFactory.getLogger("sqlServer"))) - .withExposedPorts(1433); - sqlServer.waitingFor(Wait.forLogMessage(".*(1 rows affected).*", 26)).start(); + .withEnv("ACCEPT_EULA", "Y") + .withEnv("MSSQL_PID", "Standard") + .withEnv("SA_PASSWORD", "Password!") + .withEnv("MSSQL_AGENT_ENABLED", "true") + .withLogConsumer(new Slf4jLogConsumer(LoggerFactory.getLogger("sqlServer"))) + .withExposedPorts(1433); + //sqlServer.waitingFor(Wait.forLogMessage(".*(1 rows affected).*", 50)).start(); + //sqlServer.waitingFor(Wait.forLogMessage(".*(Service Broker manager has started).*", 50)).start(); + sqlServer.start(); try (ConfigurableApplicationContext context = applicationBuilder .run("--cdc.connector=sqlserver", @@ -106,7 +117,7 @@ public class CdcSourceDatabasesIntegrationTest extends CdcTestSupport { @Test public void postgres() { - GenericContainer postgres = new GenericContainer("debezium/example-postgres:1.0") + GenericContainer postgres = new GenericContainer("debezium/example-postgres:1.3") .withEnv("POSTGRES_USER", "postgres") .withEnv("POSTGRES_PASSWORD", "postgres") .withExposedPorts(5432); diff --git a/applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattered_insert_inventory_products_106.json b/applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattened_insert_inventory_products_106.json similarity index 100% rename from applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattered_insert_inventory_products_106.json rename to applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattened_insert_inventory_products_106.json diff --git a/applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattered_update_inventory_customers.json b/applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattened_update_inventory_customers.json similarity index 100% rename from applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattered_update_inventory_customers.json rename to applications/source/cdc-debezium-source/src/test/resources/json/mysql_flattened_update_inventory_customers.json diff --git a/applications/stream-applications-core/pom.xml b/applications/stream-applications-core/pom.xml index 28473b7d..953fdccb 100644 --- a/applications/stream-applications-core/pom.xml +++ b/applications/stream-applications-core/pom.xml @@ -23,7 +23,7 @@ 1.0.0-SNAPSHOT Horsham.SR10 3.0.10.RELEASE - 1.0.0-RC1 + 1.0.0-SNAPSHOT 1.0.0-RC1 1.0.0-RC1 2.1.2.RELEASE diff --git a/functions/common/cdc-debezium-boot-starter/README.adoc b/functions/common/cdc-debezium-boot-starter/README.adoc index aa43bc12..a274e652 100644 --- a/functions/common/cdc-debezium-boot-starter/README.adoc +++ b/functions/common/cdc-debezium-boot-starter/README.adoc @@ -94,14 +94,14 @@ cdc.config.database.port=3306 # <3> cdc.schema=false # <4> -cdc.flattering.enabled=true # <5> +cdc.flattening.enabled=true # <5> ---- <1> Metadata used to identify and dispatch the events received by this cdc consumer instance. <2> Configures the CDC Source to use https://debezium.io/docs/connectors/mysql/[MySqlConnector]. (equivalent to setting `cdc.config.connector.class=io.debezium.connector.mysql.MySqlConnector`). <3> Connector specific configurations. MySQL server logical name, connect location and access credentials. <4> Do not serialize the record's schema in the output messages. -<5> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattering] feature. +<5> Enables the https://debezium.io/docs/configuration/event-flattening/[CDC Event Flattening] feature. The full list of properties: @@ -110,11 +110,11 @@ The full list of properties: //tag::configuration-properties[] $$cdc.config$$:: $$Spring pass-trough wrapper for debezium configuration properties. All properties with a 'cdc.config.' prefix are native Debezium properties. The prefix is removed, converting them into Debezium io.debezium.config.Configuration.$$ *($$Map$$, default: `$$$$`)* $$cdc.connector$$:: $$Shortcut for the cdc.config.connector.class property. Either of those can be used as long as they do not contradict with each other.$$ *($$ConnectorType$$, default: `$$$$`, possible values: `mysql`,`postgres`,`mongodb`,`oracle`,`sqlserver`)* -$$cdc.flattering.add-fields$$:: $$Comma separated list of metadata fields to add to the flattened message. The fields will be prefixed with "__" or "__[<]struct]__", depending on the specification of the struct.$$ *($$String$$, default: `$$$$`)* -$$cdc.flattering.add-headers$$:: $$Comma separated list specify a list of metadata fields to add to the header of the flattened message. The fields will be prefixed with "__" or "__[struct]__".$$ *($$String$$, default: `$$$$`)* -$$cdc.flattering.delete-handling-mode$$:: $$Options for handling deleted records: (1) none - pass the records through, (2) drop - remove the records and (3) rewrite - add a '__deleted' field to the records.$$ *($$DeleteHandlingMode$$, default: `$$$$`, possible values: `drop`,`rewrite`,`none`)* -$$cdc.flattering.drop-tombstones$$:: $$By default Debezium generates tombstone records to enable Kafka compaction on deleted records. The dropTombstones can suppress the tombstone records.$$ *($$Boolean$$, default: `$$true$$`)* -$$cdc.flattering.enabled$$:: $$Enable flattering the source record events (https://debezium.io/docs/configuration/event-flattening).$$ *($$Boolean$$, default: `$$true$$`)* +$$cdc.flattening.add-fields$$:: $$Comma separated list of metadata fields to add to the flattened message. The fields will be prefixed with "__" or "__[<]struct]__", depending on the specification of the struct.$$ *($$String$$, default: `$$$$`)* +$$cdc.flattening.add-headers$$:: $$Comma separated list specify a list of metadata fields to add to the header of the flattened message. The fields will be prefixed with "__" or "__[struct]__".$$ *($$String$$, default: `$$$$`)* +$$cdc.flattening.delete-handling-mode$$:: $$Options for handling deleted records: (1) none - pass the records through, (2) drop - remove the records and (3) rewrite - add a '__deleted' field to the records.$$ *($$DeleteHandlingMode$$, default: `$$$$`, possible values: `drop`,`rewrite`,`none`)* +$$cdc.flattening.drop-tombstones$$:: $$By default Debezium generates tombstone records to enable Kafka compaction on deleted records. The dropTombstones can suppress the tombstone records.$$ *($$Boolean$$, default: `$$true$$`)* +$$cdc.flattening.enabled$$:: $$Enable flattening the source record events (https://debezium.io/docs/configuration/event-flattening).$$ *($$Boolean$$, default: `$$true$$`)* $$cdc.name$$:: $$Unique name for this sourceConnector instance.$$ *($$String$$, default: `$$$$`)* $$cdc.offset.commit-timeout$$:: $$Maximum number of milliseconds to wait for records to flush and partition offset data to be committed to offset storage before cancelling the process and restoring the offset data to be committed in a future attempt.$$ *($$Duration$$, default: `$$5000ms$$`)* $$cdc.offset.flush-interval$$:: $$Interval at which to try committing offsets. The default is 1 minute.$$ *($$Duration$$, default: `$$60000ms$$`)* diff --git a/functions/common/cdc-debezium-boot-starter/src/main/java/org/springframework/cloud/fn/common/cdc/CdcAutoConfiguration.java b/functions/common/cdc-debezium-boot-starter/src/main/java/org/springframework/cloud/fn/common/cdc/CdcAutoConfiguration.java index e073f0cc..14e51635 100644 --- a/functions/common/cdc-debezium-boot-starter/src/main/java/org/springframework/cloud/fn/common/cdc/CdcAutoConfiguration.java +++ b/functions/common/cdc-debezium-boot-starter/src/main/java/org/springframework/cloud/fn/common/cdc/CdcAutoConfiguration.java @@ -48,10 +48,10 @@ public class CdcAutoConfiguration { @Bean public EmbeddedEngineExecutorService embeddedEngine(EmbeddedEngine.Builder embeddedEngineBuilder, - Consumer sourceRecordConsumer, Function recordFlattering) { + Consumer sourceRecordConsumer, Function recordFlattening) { EmbeddedEngine embeddedEngine = embeddedEngineBuilder - .notifying(sourceRecord -> sourceRecordConsumer.accept(recordFlattering.apply(sourceRecord))) + .notifying(sourceRecord -> sourceRecordConsumer.accept(recordFlattening.apply(sourceRecord))) .build(); return new EmbeddedEngineExecutorService(embeddedEngine) { 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 ff183b74..b37cc327 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 @@ -70,7 +70,7 @@ public class CdcBootStarterIntegrationTest { "spring.datasource.type=com.zaxxer.hikari.HikariDataSource", "cdc.name=my-sql-connector", "cdc.schema=false", - "cdc.flattering.enabled=true", + "cdc.flattening.enabled=true", "cdc.stream.header.offset=true", "cdc.connector=mysql", "cdc.config.database.user=debezium", @@ -85,8 +85,8 @@ public class CdcBootStarterIntegrationTest { public void consumerTest() { contextRunner .withPropertyValues( - "cdc.flattering.deleteHandlingMode=drop", - "cdc.flattering.dropTombstones=true") + "cdc.flattening.deleteHandlingMode=drop", + "cdc.flattening.dropTombstones=true") .run(context -> { TestCdcApplication.TestSourceRecordConsumer testConsumer = context .getBean(TestCdcApplication.TestSourceRecordConsumer.class); diff --git a/functions/common/cdc-debezium-common/pom.xml b/functions/common/cdc-debezium-common/pom.xml index efe17d63..88861be9 100644 --- a/functions/common/cdc-debezium-common/pom.xml +++ b/functions/common/cdc-debezium-common/pom.xml @@ -15,7 +15,7 @@ Change Data Capture (CDC) Debezium Common - 1.2.1.Final + 1.3.1.Final 42.2.5 5.2.1.RELEASE diff --git a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java index 434b8e7a..0f4efc4c 100644 --- a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java +++ b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java @@ -52,9 +52,9 @@ public class CdcCommonConfiguration { } @Bean - public Function recordFlattering(CdcCommonProperties properties, + public Function recordFlattening(CdcCommonProperties properties, ExtractNewRecordState extractNewRecordState) { - return sourceRecord -> properties.getFlattering().isEnabled() ? + return sourceRecord -> properties.getFlattening().isEnabled() ? (SourceRecord) extractNewRecordState.apply(sourceRecord) : sourceRecord; } @@ -62,13 +62,13 @@ public class CdcCommonConfiguration { public ExtractNewRecordState extractNewRecordState(CdcCommonProperties properties) { ExtractNewRecordState extractNewRecordState = new ExtractNewRecordState(); Map config = extractNewRecordState.config().defaultValues(); - config.put("drop.tombstones", properties.getFlattering().isDropTombstones()); - config.put("delete.handling.mode", properties.getFlattering().getDeleteHandlingMode().name()); - if (!StringUtils.isEmpty(properties.getFlattering().getAddHeaders())) { - config.put("add.headers", properties.getFlattering().getAddHeaders()); + config.put("drop.tombstones", properties.getFlattening().isDropTombstones()); + config.put("delete.handling.mode", properties.getFlattening().getDeleteHandlingMode().name()); + if (!StringUtils.isEmpty(properties.getFlattening().getAddHeaders())) { + config.put("add.headers", properties.getFlattening().getAddHeaders()); } - if (!StringUtils.isEmpty(properties.getFlattering().getAddFields())) { - config.put("add.fields", properties.getFlattering().getAddFields()); + if (!StringUtils.isEmpty(properties.getFlattening().getAddFields())) { + config.put("add.fields", properties.getFlattening().getAddFields()); } extractNewRecordState.configure(config); diff --git a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java index 7756c75b..93e146d1 100644 --- a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java +++ b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java @@ -56,9 +56,9 @@ public class CdcCommonProperties { private boolean schema = false; /** - * Event Flattering (https://debezium.io/docs/configuration/event-flattening). + * Event Flattening (https://debezium.io/docs/configuration/event-flattening). */ - private final Flattering flattering = new Flattering(); + private final Flattening flattening = new Flattening(); /** * Spring pass-trough wrapper for debezium configuration properties. @@ -79,8 +79,8 @@ public class CdcCommonProperties { return offset; } - public Flattering getFlattering() { - return flattering; + public Flattening getFlattening() { + return flattening; } public Map getConfig() { @@ -238,10 +238,10 @@ public class CdcCommonProperties { * https://debezium.io/documentation/reference/0.10/configuration/event-flattening.html . * https://debezium.io/documentation/reference/0.10/configuration/event-flattening.html#configuration_options */ - public static class Flattering { + public static class Flattening { /** - * Enable flattering the source record events (https://debezium.io/docs/configuration/event-flattening). + * Enable flattening the source record events (https://debezium.io/docs/configuration/event-flattening). */ private boolean enabled = true; 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 03fec9f9..624653be 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 @@ -218,7 +218,6 @@ public final class EmbeddedEngine implements DebeziumEngine { public static final class BuilderImpl implements Builder { private OffsetBackingStore offsetBackingStore; - private SourceConnector sourceConnector; private Configuration config; private DebeziumEngine.ChangeConsumer handler; private ClassLoader classLoader; @@ -227,12 +226,6 @@ public final class EmbeddedEngine implements DebeziumEngine { private DebeziumEngine.ConnectorCallback connectorCallback; private OffsetCommitPolicy offsetCommitPolicy = null; - @Override - public Builder sourceConnector(SourceConnector sourceConnector) { - this.sourceConnector = sourceConnector; - return this; - } - @Override public Builder offsetBackingStore(OffsetBackingStore offsetBackingStore) { this.offsetBackingStore = offsetBackingStore; @@ -315,8 +308,7 @@ public final class EmbeddedEngine implements DebeziumEngine { Objects.requireNonNull(config, "A connector configuration must be specified."); Objects.requireNonNull(handler, "A connector consumer or changeHandler must be specified."); return new EmbeddedEngine(config, classLoader, clock, - handler, completionCallback, connectorCallback, offsetCommitPolicy, - sourceConnector, offsetBackingStore); + handler, completionCallback, connectorCallback, offsetCommitPolicy, offsetBackingStore); } // backward compatibility methods @@ -541,8 +533,6 @@ public final class EmbeddedEngine implements DebeziumEngine { @Override Builder using(OffsetCommitPolicy policy); - Builder sourceConnector(SourceConnector sourceConnector); - Builder offsetBackingStore(OffsetBackingStore offsetBackingStore); @Override @@ -576,7 +566,6 @@ public final class EmbeddedEngine implements DebeziumEngine { private long recordsSinceLastCommit = 0; private long timeOfLastCommitMillis = 0; private OffsetCommitPolicy offsetCommitPolicy; - private SourceConnector connector; private OffsetBackingStore offsetStore; private SourceTask task; @@ -584,8 +573,7 @@ public final class EmbeddedEngine implements DebeziumEngine { private EmbeddedEngine(Configuration config, ClassLoader classLoader, Clock clock, DebeziumEngine.ChangeConsumer handler, DebeziumEngine.CompletionCallback completionCallback, DebeziumEngine.ConnectorCallback connectorCallback, - OffsetCommitPolicy offsetCommitPolicy, SourceConnector sourceConnector, OffsetBackingStore offsetStore) { - this.connector = sourceConnector; + OffsetCommitPolicy offsetCommitPolicy, OffsetBackingStore offsetStore) { this.offsetStore = offsetStore; this.config = config; this.handler = handler; diff --git a/functions/common/cdc-debezium-common/src/main/resources/application.properties b/functions/common/cdc-debezium-common/src/main/resources/application.properties index 8fb76f9d..6bc13f69 100644 --- a/functions/common/cdc-debezium-common/src/main/resources/application.properties +++ b/functions/common/cdc-debezium-common/src/main/resources/application.properties @@ -1,2 +1,2 @@ spring.autoconfigure.exclude=org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration -spring.cloud.stream.kafka.default.producer.messageKeyExpression=headers['cdc_key'] +spring.cloud.stream.kafka.default.producer.messageKeyExpression=headers['cdc_key'].bytes diff --git a/functions/supplier/cdc-debezium-supplier/pom.xml b/functions/supplier/cdc-debezium-supplier/pom.xml index c1cfdb75..95a97976 100644 --- a/functions/supplier/cdc-debezium-supplier/pom.xml +++ b/functions/supplier/cdc-debezium-supplier/pom.xml @@ -14,7 +14,7 @@ CDC Debezium Suppliers - 1.2.1.Final + 1.3.1.Final 8.0.13 diff --git a/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java b/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java index 9055e3fa..ba546f23 100644 --- a/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java +++ b/functions/supplier/cdc-debezium-supplier/src/main/java/org/springframework/cloud/fn/supplier/cdc/CdcSupplierConfiguration.java @@ -111,13 +111,13 @@ public class CdcSupplierConfiguration implements BeanClassLoaderAware { public EmbeddedEngineExecutorService embeddedEngineExecutorService( EmbeddedEngine.Builder embeddedEngineBuilder, Function valueSerializer, Function keySerializer, - Function recordFlattering, + Function recordFlattening, ObjectMapper mapper, CdcSupplierProperties cdcStreamingEngineProperties) { FluxSink> sink = emitterProcessor.sink(); Consumer messageConsumer = sourceRecord -> { - // When cdc.flattering.deleteHandlingMode=none and cdc.flattering.dropTombstones=false + // When cdc.flattening.deleteHandlingMode=none and cdc.flattening.dropTombstones=false // then on deletion event an additional sourceRecord is sent with value Null. // Here we filter out such condition. if (sourceRecord == null) { @@ -130,7 +130,7 @@ public class CdcSupplierConfiguration implements BeanClassLoaderAware { // When the tombstone event is enabled, Debezium serializes the payload to null (e.g. // empty payload) // while the metadata information is carried through the headers (cdc_key). - // Note: Event for none flattered responses, when the cdc.config.tombstones.on.delete=true + // Note: Event for none flattened responses, when the cdc.config.tombstones.on.delete=true // (default), // tombstones are generate by Debezium and handled by the code below. if (cdcJsonPayload == null) { @@ -151,7 +151,8 @@ public class CdcSupplierConfiguration implements BeanClassLoaderAware { MessageBuilder messageBuilder = MessageBuilder .withPayload(cdcJsonPayload) - .setHeader("cdc_key", new String(key)) +// .setHeader("cdc_key", new String(key)) + .setHeader("cdc_key", key) .setHeader("cdc_topic", sourceRecord.topic()) .setHeader(MessageHeaders.CONTENT_TYPE, (cdcJsonPayload.equals(this.kafkaNull)) ? MimeTypeUtils.TEXT_PLAIN_VALUE @@ -182,7 +183,7 @@ public class CdcSupplierConfiguration implements BeanClassLoaderAware { }; EmbeddedEngine engine = embeddedEngineBuilder - .notifying(record -> messageConsumer.accept(recordFlattering.apply(record))) + .notifying(record -> messageConsumer.accept(recordFlattening.apply(record))) .build(); return new EmbeddedEngineExecutorService(engine); diff --git a/functions/supplier/cdc-debezium-supplier/src/main/resources/application.properties b/functions/supplier/cdc-debezium-supplier/src/main/resources/application.properties index 8fb76f9d..6bc13f69 100644 --- a/functions/supplier/cdc-debezium-supplier/src/main/resources/application.properties +++ b/functions/supplier/cdc-debezium-supplier/src/main/resources/application.properties @@ -1,2 +1,2 @@ spring.autoconfigure.exclude=org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration -spring.cloud.stream.kafka.default.producer.messageKeyExpression=headers['cdc_key'] +spring.cloud.stream.kafka.default.producer.messageKeyExpression=headers['cdc_key'].bytes