Upgrade Debezium to 2.3.0.Final
* address some review comments * minor pom improvements * undo Debezium BOM
This commit is contained in:
@@ -1,11 +1,11 @@
|
||||
//tag::ref-doc[]
|
||||
= Debezium Source
|
||||
|
||||
https://debezium.io/documentation/reference/2.2/development/engine.html[Debezium Engine] based https://en.wikipedia.org/wiki/Change_data_capture[Change Data Capture] (CDC) source.
|
||||
https://debezium.io/documentation/reference/development/engine.html[Debezium Engine] based https://en.wikipedia.org/wiki/Change_data_capture[Change Data Capture] (CDC) source.
|
||||
The `Debezium Source` allows *capturing* database change events and *streaming* those over different message binders such `Apache Kafka`, `RabbitMQ` and all Spring Cloud Stream supporter brokers.
|
||||
|
||||
NOTE: This source can be used with *any* Spring Cloud Stream message binder.
|
||||
It is not restricted nor depended on the Kafka Connect framework. Though this approach is flexible it comes with certain https://debezium.io/documentation/reference/2.2/development/engine.html#_handling_failures[limitations].
|
||||
It is not restricted nor depended on the Kafka Connect framework. Though this approach is flexible it comes with certain https://debezium.io/documentation/reference/development/engine.html#_handling_failures[limitations].
|
||||
|
||||
All Debezium configuration properties are supported.
|
||||
Just precede any Debezium properties with the `debezium.properties.` prefix.
|
||||
@@ -13,7 +13,7 @@ For example to set the Debezium's `connector.class` property use the `debezium.p
|
||||
|
||||
== Database Support
|
||||
|
||||
The `Debezium Source` currently supports CDC for multiple datastores: https://debezium.io/documentation/reference/2.2/connectors/mysql.html[MySQL], https://debezium.io/documentation/reference/2.2/connectors/postgresql.html[PostgreSQL], https://debezium.io/documentation/reference/2.2/connectors/mongodb.html[MongoDB], https://debezium.io/documentation/reference/2.2/connectors/oracle.html[Oracle], https://debezium.io/documentation/reference/2.2/connectors/sqlserver.html[SQL Server], https://debezium.io/documentation/reference/2.2/connectors/db2.html[Db2], https://debezium.io/documentation/reference/2.2/connectors/vitess.html[Vitess] and https://debezium.io/documentation/reference/2.2/connectors/spanner.html[Spanner] databases.
|
||||
The `Debezium Source` currently supports CDC for multiple datastores: https://debezium.io/documentation/reference/connectors/mysql.html[MySQL], https://debezium.io/documentation/reference/connectors/postgresql.html[PostgreSQL], https://debezium.io/documentation/reference/connectors/mongodb.html[MongoDB], https://debezium.io/documentation/reference/connectors/oracle.html[Oracle], https://debezium.io/documentation/reference/connectors/sqlserver.html[SQL Server], https://debezium.io/documentation/reference/connectors/db2.html[Db2], https://debezium.io/documentation/reference/connectors/vitess.html[Vitess] and https://debezium.io/documentation/reference/connectors/spanner.html[Spanner] databases.
|
||||
|
||||
== Options
|
||||
|
||||
@@ -55,7 +55,7 @@ Using the https://debezium.io/documentation/reference/stable/transformations/eve
|
||||
|
||||
When a Debezium source runs, it reads information from the source and periodically records `offsets` that define how much of that information it has processed.
|
||||
Should the source be restarted, it will use the last recorded offset to know where in the source information it should resume reading.
|
||||
Out of the box, the following https://debezium.io/documentation/reference/2.2/development/engine.html#engine-properties[offset storage configuration] options are provided:
|
||||
Out of the box, the following https://debezium.io/documentation/reference/development/engine.html#engine-properties[offset storage configuration] options are provided:
|
||||
|
||||
- In-Memory
|
||||
|
||||
@@ -104,32 +104,32 @@ Those properties can be used by prefixing them by the `debezium.properties.` pre
|
||||
|===
|
||||
| Connector | Connector properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mysql.html[MySQL]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mysql.html#mysql-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/mysql.html[MySQL]
|
||||
|https://debezium.io/documentation/reference/connectors/mysql.html#mysql-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mongodb.html[MongoDB]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mongodb.html#mongodb-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/mongodb.html[MongoDB]
|
||||
|https://debezium.io/documentation/reference/connectors/mongodb.html#mongodb-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/postgresql.html[PostgreSQL]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/postgresql.html#postgresql-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/postgresql.html[PostgreSQL]
|
||||
|https://debezium.io/documentation/reference/connectors/postgresql.html#postgresql-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/oracle.html[Oracle]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/oracle.html#oracle-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/oracle.html[Oracle]
|
||||
|https://debezium.io/documentation/reference/connectors/oracle.html#oracle-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/sqlserver.html[SQL Server]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/sqlserver.html#sqlserver-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/sqlserver.html[SQL Server]
|
||||
|https://debezium.io/documentation/reference/connectors/sqlserver.html#sqlserver-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/db2.html[DB2]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/db2.html#db2-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/db2.html[DB2]
|
||||
|https://debezium.io/documentation/reference/connectors/db2.html#db2-connector-properties
|
||||
|
||||
// |https://debezium.io/documentation/reference/2.2/connectors/cassandra.html[Cassandra]
|
||||
// |https://debezium.io/documentation/reference/2.2/connectors/cassandra.html#cassandra-connector-properties
|
||||
// |https://debezium.io/documentation/reference/connectors/cassandra.html[Cassandra]
|
||||
// |https://debezium.io/documentation/reference/connectors/cassandra.html#cassandra-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/vitess.html[Vitess]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/vitess.html#vitess-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/vitess.html[Vitess]
|
||||
|https://debezium.io/documentation/reference/connectors/vitess.html#vitess-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/spanner.html[Spanner]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/spanner.html#spanner-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/spanner.html[Spanner]
|
||||
|https://debezium.io/documentation/reference/connectors/spanner.html#spanner-connector-properties
|
||||
|
||||
|===
|
||||
|
||||
@@ -145,7 +145,7 @@ Instructions below explains how to run pre-configured test databases form Docker
|
||||
Start the `debezium/example-mysql` in a docker:
|
||||
[source, bash]
|
||||
----
|
||||
docker run -it --rm --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=debezium -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw debezium/example-mysql:2.2.0.Final
|
||||
docker run -it --rm --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=debezium -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw debezium/example-mysql:2.3.0.Final
|
||||
----
|
||||
|
||||
[TIP]
|
||||
@@ -192,7 +192,7 @@ debezium.properties.offset.storage=org.apache.kafka.connect.storage.MemoryOffset
|
||||
<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 `ChangeEvent` message.
|
||||
<5> Enables the https://debezium.io/documentation/reference/2.2/transformations/event-flattening.html[Change Event Flattening].
|
||||
<5> Enables the https://debezium.io/documentation/reference/transformations/event-flattening.html[Change Event Flattening].
|
||||
<6> Source state to preserver between multiple starts.
|
||||
|
||||
You can run also the `DebeziumDatabasesIntegrationTest#mysql()` using this mysql configuration.
|
||||
@@ -205,7 +205,7 @@ NOTE: Disable the mysql GenericContainer test initialization code.
|
||||
Start a pre-configured postgres server from the `debezium/example-postgres:1.0` Docker image:
|
||||
[source, bash]
|
||||
----
|
||||
docker run -it --rm --name postgres -p 5432:5432 -e POSTGRES_USER=postgres -e POSTGRES_PASSWORD=postgres debezium/example-postgres:2.2.0.Final
|
||||
docker run -it --rm --name postgres -p 5432:5432 -e POSTGRES_USER=postgres -e POSTGRES_PASSWORD=postgres debezium/example-postgres:2.3.0.Final
|
||||
----
|
||||
|
||||
You can connect to this server like this:
|
||||
@@ -257,10 +257,10 @@ NOTE: Disable the postgres GenericContainer test initialization code.
|
||||
|
||||
=== MongoDB
|
||||
|
||||
Start a pre-configured mongodb from the `debezium/example-mongodb:2.2.0.Final` container image:
|
||||
Start a pre-configured mongodb from the `debezium/example-mongodb:2.3.0.Final` container image:
|
||||
[source, bash]
|
||||
----
|
||||
docker run -it --rm --name mongodb -p 27017:27017 -e MONGODB_USER=debezium -e MONGODB_PASSWORD=dbz debezium/example-mongodb:2.2.0.Final
|
||||
docker run -it --rm --name mongodb -p 27017:27017 -e MONGODB_USER=debezium -e MONGODB_PASSWORD=dbz debezium/example-mongodb:2.3.0.Final
|
||||
----
|
||||
|
||||
Initialize the inventory collections
|
||||
|
||||
@@ -115,7 +115,7 @@
|
||||
<dependency>
|
||||
<groupId>io.debezium</groupId>
|
||||
<artifactId>debezium-testing-testcontainers</artifactId>
|
||||
<version>2.2.0.Final</version>
|
||||
<version>2.3.0.Final</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
|
||||
@@ -107,7 +107,7 @@ public class DebeziumDatabasesIntegrationTest {
|
||||
|
||||
assertThat(messages).isNotNull();
|
||||
// Message size should correspond to the number of insert statements in:
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mysql/2.2/inventory.sql
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mysql/2.3/inventory.sql
|
||||
assertThat(messages).hasSizeGreaterThanOrEqualTo(52);
|
||||
}
|
||||
mySQL.stop();
|
||||
@@ -144,7 +144,7 @@ public class DebeziumDatabasesIntegrationTest {
|
||||
allMessages.addAll(messageChunk);
|
||||
}
|
||||
// Message size should correspond to the number of insert statements in the sample inventor DB:
|
||||
// https://github.com/debezium/container-images/blob/main/examples/postgres/2.2/inventory.sql
|
||||
// https://github.com/debezium/container-images/blob/main/examples/postgres/2.3/inventory.sql
|
||||
return allMessages.size() == 29; // Inventory DB entries
|
||||
});
|
||||
}
|
||||
@@ -234,7 +234,7 @@ public class DebeziumDatabasesIntegrationTest {
|
||||
List<Message<?>> messages = DebeziumTestUtils.receiveAll(outputDestination);
|
||||
assertThat(messages).isNotNull();
|
||||
// Number of entries should match the entries inserted by:
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mongodb/2.2/init-inventory.sh
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mongodb/2.3/init-inventory.sh
|
||||
assertThat(messages).hasSize(666);
|
||||
}
|
||||
mongodb.stop();
|
||||
|
||||
@@ -34,7 +34,7 @@ public final class DebeziumTestUtils {
|
||||
|
||||
public static final String DATABASE_NAME = "inventory";
|
||||
public static final String BINDING_NAME = "debeziumSupplier-out-0";
|
||||
public static final String IMAGE_TAG = "2.2.0.Final";
|
||||
public static final String IMAGE_TAG = "2.3.0.Final";
|
||||
public static final String DEBEZIUM_EXAMPLE_MYSQL_IMAGE = "debezium/example-mysql:" + IMAGE_TAG;
|
||||
public static final String DEBEZIUM_EXAMPLE_POSTGRES_IMAGE = "debezium/example-postgres:" + IMAGE_TAG;
|
||||
public static final String DEBEZIUM_EXAMPLE_MONGODB_IMAGE = "debezium/example-mongodb:" + IMAGE_TAG;
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
= Debezium Auto-Configuration
|
||||
|
||||
This module provides a generic https://debezium.io/documentation/reference/2.2/development/engine.html[DebeziumEngine.Builder] auto-configuration that can be reused and composed in other applications.
|
||||
This module provides a generic https://debezium.io/documentation/reference/development/engine.html[DebeziumEngine.Builder] auto-configuration that can be reused and composed in other applications.
|
||||
|
||||
IMPORTANT: The `DebeziumEngine` does not required Kafka or Kafka Connect as it runs embedded inside your application.
|
||||
This approach though comes with some delivery guarantee limitations as explained https://debezium.io/documentation/reference/2.2/development/engine.html#%5Fhandling_failures[here].
|
||||
This approach though comes with some delivery guarantee limitations as explained https://debezium.io/documentation/reference/development/engine.html#%5Fhandling_failures[here].
|
||||
|
||||
The `Debezium Engine` is a https://en.wikipedia.org/wiki/Change_data_capture[Change Data Capture] (CDC) utility, that allows *capturing* database change events and process them with custom `java.util.Consumer` or `io.debezium.engine.ChangeConsumer` event handler implementations.
|
||||
|
||||
@@ -36,7 +36,7 @@ compile "org.springframework.cloud.fn:debezium-autoconfigure:{project-version}"
|
||||
----
|
||||
====
|
||||
|
||||
and include the https://debezium.io/documentation/reference/2.2/connectors/index.html[debezium connector] dependency for the selected Database.
|
||||
and include the https://debezium.io/documentation/reference/connectors/index.html[debezium connector] dependency for the selected Database.
|
||||
For example the postgres debezium connector dependency looks like this:
|
||||
|
||||
====
|
||||
@@ -178,29 +178,29 @@ The table below lists all available Debezium properties for each connecter.
|
||||
|===
|
||||
| Connector | Connector properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mysql.html[MySQL]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mysql.html#mysql-connector-properties
|
||||
|https://debezium.io/documentation/reference/reference/connectors/mysql.html[MySQL]
|
||||
|https://debezium.io/documentation/reference/connectors/mysql.html#mysql-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mongodb.html[MongoDB]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/mongodb.html#mongodb-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/mongodb.html[MongoDB]
|
||||
|https://debezium.io/documentation/reference/connectors/mongodb.html#mongodb-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/postgresql.html[PostgreSQL]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/postgresql.html#postgresql-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/postgresql.html[PostgreSQL]
|
||||
|https://debezium.io/documentation/reference/connectors/postgresql.html#postgresql-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/oracle.html[Oracle]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/oracle.html#oracle-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/oracle.html[Oracle]
|
||||
|https://debezium.io/documentation/reference/connectors/oracle.html#oracle-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/sqlserver.html[SQL Server]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/sqlserver.html#sqlserver-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/sqlserver.html[SQL Server]
|
||||
|https://debezium.io/documentation/reference/connectors/sqlserver.html#sqlserver-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/db2.html[DB2]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/db2.html#db2-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/db2.html[DB2]
|
||||
|https://debezium.io/documentation/reference/connectors/db2.html#db2-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/vitess.html[Vitess]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/vitess.html#vitess-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/vitess.html[Vitess]
|
||||
|https://debezium.io/documentation/reference/connectors/vitess.html#vitess-connector-properties
|
||||
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/spanner.html[Spanner]
|
||||
|https://debezium.io/documentation/reference/2.2/connectors/spanner.html#spanner-connector-properties
|
||||
|https://debezium.io/documentation/reference/connectors/spanner.html[Spanner]
|
||||
|https://debezium.io/documentation/reference/connectors/spanner.html#spanner-connector-properties
|
||||
|
||||
|===
|
||||
|
||||
@@ -287,7 +287,7 @@ Follow the https://debezium.io/documentation/reference/stable/transformations/ev
|
||||
|
||||
When a Debezium source runs, it reads information from the source and periodically records `offsets` that define how much of that information it has processed.
|
||||
Should the source be restarted, it will use the last recorded offset to know where in the source information it should resume reading.
|
||||
Out of the box, the following https://debezium.io/documentation/reference/2.2/development/engine.html#engine-properties[offset storage configuration] options are provided:
|
||||
Out of the box, the following https://debezium.io/documentation/reference/development/engine.html#engine-properties[offset storage configuration] options are provided:
|
||||
|
||||
==== In-Memory
|
||||
|
||||
|
||||
@@ -15,7 +15,7 @@
|
||||
<description>Debezium Spring Boot auto-configuration</description>
|
||||
|
||||
<properties>
|
||||
<version.debezium>2.3.0.CR1</version.debezium>
|
||||
<version.debezium>2.3.0.Final</version.debezium>
|
||||
<apicurio.version>2.4.2.Final</apicurio.version>
|
||||
</properties>
|
||||
|
||||
@@ -26,12 +26,16 @@
|
||||
<version>${version.debezium}</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<artifactId>slf4j-reload4j</artifactId>
|
||||
<groupId>org.slf4j</groupId>
|
||||
<artifactId>slf4j-reload4j</artifactId>
|
||||
</exclusion>
|
||||
<exclusion>
|
||||
<artifactId>slf4j-log4j12</artifactId>
|
||||
<groupId>org.slf4j</groupId>
|
||||
<artifactId>slf4j-log4j12</artifactId>
|
||||
</exclusion>
|
||||
<exclusion>
|
||||
<groupId>org.slf4j</groupId>
|
||||
<artifactId>slf4j-api</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
@@ -46,6 +50,12 @@
|
||||
<groupId>io.apicurio</groupId>
|
||||
<artifactId>apicurio-registry-client</artifactId>
|
||||
<version>${apicurio.version}</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.slf4j</groupId>
|
||||
<artifactId>slf4j-api</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
|
||||
<!-- TEST -->
|
||||
@@ -72,6 +82,12 @@
|
||||
<artifactId>spring-jdbc</artifactId>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>ch.qos.logback</groupId>
|
||||
<artifactId>logback-classic</artifactId>
|
||||
<version>1.4.8</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
<build>
|
||||
<plugins>
|
||||
|
||||
@@ -22,7 +22,6 @@ import java.util.Objects;
|
||||
|
||||
import io.debezium.engine.ChangeEvent;
|
||||
import io.debezium.engine.DebeziumEngine;
|
||||
import io.debezium.engine.DebeziumEngine.Builder;
|
||||
import io.debezium.engine.DebeziumEngine.CompletionCallback;
|
||||
import io.debezium.engine.DebeziumEngine.ConnectorCallback;
|
||||
import io.debezium.engine.format.KeyValueHeaderChangeEventFormat;
|
||||
@@ -123,6 +122,7 @@ public class DebeziumEngineBuilderAutoConfiguration {
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public DebeziumEngine.Builder<ChangeEvent<byte[], byte[]>> debeziumEngineBuilder(
|
||||
OffsetCommitPolicy offsetCommitPolicy, CompletionCallback completionCallback,
|
||||
ConnectorCallback connectorCallback, DebeziumProperties properties, Clock debeziumClock) {
|
||||
|
||||
@@ -51,7 +51,6 @@ import org.springframework.context.annotation.Primary;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.test.jdbc.JdbcTestUtils;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.awaitility.Awaitility.await;
|
||||
|
||||
/**
|
||||
@@ -66,7 +65,7 @@ public class DebeziumEngineBuilderAutoConfigurationIntegrationTest {
|
||||
private static final Log logger = LogFactory.getLog(DebeziumEngineBuilderAutoConfigurationIntegrationTest.class);
|
||||
|
||||
private static final String DATABASE_NAME = "inventory";
|
||||
public static final String IMAGE_TAG = "2.2.0.Final";
|
||||
public static final String IMAGE_TAG = "2.3.0.Final";
|
||||
public static final String DEBEZIUM_EXAMPLE_MYSQL_IMAGE = "debezium/example-mysql:" + IMAGE_TAG;
|
||||
|
||||
@TempDir
|
||||
@@ -136,8 +135,7 @@ public class DebeziumEngineBuilderAutoConfigurationIntegrationTest {
|
||||
"VALUES('Test666', 'Test666', 'Test666@spring.org')");
|
||||
JdbcTestUtils.deleteFromTableWhere(jdbcTemplate, "customers", "first_name = ?", "Test666");
|
||||
|
||||
await().atMost(Duration.ofSeconds(30))
|
||||
.untilAsserted(() -> assertThat(testConsumer.recordList).hasSizeGreaterThanOrEqualTo(52));
|
||||
await().atMost(Duration.ofSeconds(30)).until(() -> (testConsumer.recordList.size() >= 52));
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -15,7 +15,7 @@
|
||||
<description>Debezium Supplier</description>
|
||||
|
||||
<properties>
|
||||
<version.debezium>2.3.0.CR1</version.debezium>
|
||||
<version.debezium>2.3.0.Final</version.debezium>
|
||||
</properties>
|
||||
|
||||
<dependencies>
|
||||
|
||||
@@ -40,7 +40,7 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = {
|
||||
"spring.cloud.function.definition=debeziumSupplier",
|
||||
|
||||
// https://debezium.io/documentation/reference/2.2/transformations/event-flattening.html
|
||||
// https://debezium.io/documentation/reference/transformations/event-flattening.html
|
||||
"debezium.properties.transforms=unwrap",
|
||||
"debezium.properties.transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState",
|
||||
"debezium.properties.transforms.unwrap.drop.tombstones=true",
|
||||
@@ -73,7 +73,7 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
@Testcontainers
|
||||
public class DebeziumSupplierIntegrationTest {
|
||||
|
||||
public static final String IMAGE_TAG = "2.2.0.Final";
|
||||
public static final String IMAGE_TAG = "2.3.0.Final";
|
||||
public static final String DEBEZIUM_EXAMPLE_MYSQL_IMAGE = "debezium/example-mysql:" + IMAGE_TAG;
|
||||
|
||||
@Container
|
||||
@@ -106,7 +106,7 @@ public class DebeziumSupplierIntegrationTest {
|
||||
Flux<Message<?>> messageFlux = this.debeziumSupplier.get();
|
||||
|
||||
// Message size should correspond to the number of insert statements in:
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mysql/2.2/inventory.sql
|
||||
// https://github.com/debezium/container-images/blob/main/examples/mysql/2.3/inventory.sql
|
||||
// filtered by Customers and Addresses table.
|
||||
StepVerifier.create(messageFlux)
|
||||
.expectNextCount(16) // Skip the DDL transaction logs.
|
||||
|
||||
Reference in New Issue
Block a user