集成

使用 Testcontainers 进行集成测试

师成师成· 更新于 2026-09-28· 阅读 12 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

使用 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 的变更事件主题中读取两条记录,并断言其属性

评论

登录后参与评论

正在加载评论…