集成
使用 Testcontainers 进行集成测试
登录后可跨设备保存划线和私人笔记登录
使用 Testcontainers 进行集成测试
概述
在使用 Debezium 搭建变更数据捕获(CDC)管道时,最好同时建立一些自动化测试,以确保:
- 源数据库已完成相应设置,变更能够从中流出
- 连接器的配置正确无误
Testcontainers 的 Debezium 扩展旨在简化这类测试:它通过 Linux 容器运行所有必需的基础设施(Apache Kafka、Kafka Connect 等),并使其能够被基于 Java 的测试轻松访问。
该扩展尽可能采用合理的默认值(例如,连接器的数据库凭据可以直接从已配置的数据库容器中获取),从而让你能够专注于测试的核心逻辑。
快速开始
要使用 Debezium 的 Testcontainers 集成功能,请在项目中添加以下依赖:
<dependency>
<groupId>io.debezium</groupId>
<artifactId>debezium-testing-testcontainers</artifactId>
<version>3.6.3.Final</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>kafka</artifactId>
<scope>test</scope>
</dependency>
<!-- Add the TC dependency matching your database -->
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>postgresql</artifactId>
<scope>test</scope>
</dependency>根据你的测试策略,你可能还需要数据库的 JDBC 驱动以及 Apache Kafka 的客户端,以便插入测试数据并断言 Kafka 中相应的变更事件。
测试环境搭建
为 Debezium 连接器配置编写集成测试时,你还需要搭建 Apache Kafka 以及作为变更事件来源的数据库。为此可以使用现有的 Testcontainers 对 Apache Kafka 和 数据库 的支持。
结合 Debezium 的 DebeziumContainer 类,典型的环境搭建大致如下:
public class DebeziumContainerTest {
private static Network network = Network.newNetwork(); (1)
private static KafkaContainer kafkaContainer = new KafkaContainer()
.withNetwork(network); (2)
public static PostgreSQLContainer<?> postgresContainer =
new PostgreSQLContainer<>(
DockerImageName.parse("quay.io/debezium/postgres:15")
.asCompatibleSubstituteFor("postgres"))
.withNetwork(network)
.withNetworkAliases("postgres"); (3)
public static DebeziumContainer debeziumContainer =
new DebeziumContainer("quay.io/debezium/connect:3.6.3.Final")
.withNetwork(network)
.withKafka(kafkaContainer)
.dependsOn(kafkaContainer); (4)
@BeforeClass
public static void startContainers() { (5)
Startables.deepStart(Stream.of(
kafkaContainer, postgresContainer, debeziumContainer))
.join();
}
}| 1 | 定义一个供所有服务共用的 Docker 网络 |
|---|---|
| 2 | 为 Apache Kafka 准备一个容器 |
| 3 | 为 Postgres 15 准备一个容器(使用 Debezium 的 Postgres 容器镜像) |
| 4 | 为 Kafka Connect 准备一个搭载 Debezium 3.6.3.Final 的容器 |
| 5 | 启动全部三个容器 |
测试实现
在声明了所有必需的容器之后,现在可以注册一个 Debezium Postgres 连接器实例,向 Postgres 插入一些测试数据,然后使用 Apache Kafka 客户端从相应的 Kafka 主题中读取预期的变更事件记录:
@Test
public void canRegisterPostgreSqlConnector() throws Exception {
try (Connection connection = getConnection(postgresContainer);
Statement statement = connection.createStatement();
KafkaConsumer<String, String> consumer = getConsumer(
kafkaContainer)) {
statement.execute("create schema todo"); (1)
statement.execute("create table todo.Todo (id int8 not null, " +
"title varchar(255), primary key (id))");
statement.execute("alter table todo.Todo replica identity full");
statement.execute("insert into todo.Todo values (1, " +
"'Learn CDC')");
statement.execute("insert into todo.Todo values (2, " +
"'Learn Debezium')");
ConnectorConfiguration connector = ConnectorConfiguration
.forJdbcContainer(postgresContainer)
.with("topic.prefix", "dbserver1");
debeziumContainer.registerConnector("my-connector",
connector); (2)
consumer.subscribe(Arrays.asList("dbserver1.todo.todo"));
List<ConsumerRecord<String, String>> changeEvents =
drain(consumer, 2); (3)
assertThat(JsonPath.<Integer> read(changeEvents.get(0).key(),
"$.id")).isEqualTo(1);
assertThat(JsonPath.<String> read(changeEvents.get(0).value(),
"$.op")).isEqualTo("r");
assertThat(JsonPath.<String> read(changeEvents.get(0).value(),
"$.after.title")).isEqualTo("Learn CDC");
assertThat(JsonPath.<Integer> read(changeEvents.get(1).key(),
"$.id")).isEqualTo(2);
assertThat(JsonPath.<String> read(changeEvents.get(1).value(),
"$.op")).isEqualTo("r");
assertThat(JsonPath.<String> read(changeEvents.get(1).value(),
"$.after.title")).isEqualTo("Learn Debezium");
consumer.unsubscribe();
}
}
// Helper methods below
private Connection getConnection(
PostgreSQLContainer<?> postgresContainer)
throws SQLException {
return DriverManager.getConnection(postgresContainer.getJdbcUrl(),
postgresContainer.getUsername(),
postgresContainer.getPassword());
}
private KafkaConsumer<String, String> getConsumer(
KafkaContainer kafkaContainer) {
return new KafkaConsumer<>(
Map.of(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
kafkaContainer.getBootstrapServers(),
ConsumerConfig.GROUP_ID_CONFIG,
"tc-" + UUID.randomUUID(),
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"earliest"),
new StringDeserializer(),
new StringDeserializer());
}
private List<ConsumerRecord<String, String>> drain(
KafkaConsumer<String, String> consumer,
int expectedRecordCount) {
List<ConsumerRecord<String, String>> allRecords = new ArrayList<>();
Unreliables.retryUntilTrue(10, TimeUnit.SECONDS, () -> {
consumer.poll(Duration.ofMillis(50))
.iterator()
.forEachRemaining(allRecords::add);
return allRecords.size() == expectedRecordCount;
});
return allRecords;
}| 1 | 在 Postgres 数据库中创建一张表并插入两条记录 |
|---|---|
| 2 | 注册一个 Debezium Postgres 连接器实例;连接器类型以及数据库主机、数据库名称、用户等属性均从数据库容器中获取 |
| 3 | 从 Kafka 的变更事件主题中读取两条记录,并断言其属性 |
评论
登录后参与评论
正在加载评论…
KnowForge