使用 Testcontainers 的 Pulsar 模块:在 JUnit 测试中一键启动 Apache Pulsar 容器
使用 Testcontainers 的 Pulsar 模块在 JUnit 测试中一键启动 Apache Pulsar 容器【免费下载链接】testcontainers-javaTestcontainers is a Java library that supports JUnit tests, providing lightweight, throwaway instances of common databases, Selenium web browsers, or anything else that can run in a Docker container.项目地址: https://gitcode.com/GitHub_Trending/te/testcontainers-javaTestcontainers 的 Pulsar 模块testcontainers-pulsar可以在 JUnit 测试中自动创建 Apache Pulsar 的 Docker 容器无需任何外部服务即可验证生产者、消费者、管理 API、Pulsar IO 乃至事务等完整能力。读完本文你将掌握如何用几行代码启动 Pulsar standalone 容器并获取 broker/HTTP 地址如何通过PULSAR_PREFIX_环境变量注入任意 Pulsar 配置以及如何按需开启 Functions Worker 与事务支持并理解容器底层启动命令与就绪等待策略的实现原理。模块概述与设计目标Pulsar 模块基于官方 Apache Pulsar Docker 镜像apachepulsar/pulsar与apachepulsar/pulsar-all由容器封装了 Pulsar 的standalone 模式单容器内同时承载 broker 与 BookKeeper适合测试场景。官方推荐在使用前阅读 Apache Pulsar 官方 Docker 入门指南 以了解镜像本身的行为。模块核心类是 org.testcontainers.pulsar.PulsarContainer它继承自GenericContainerPulsarContainer公开了两个固定端口常量常量端口用途BROKER_PORT6650Pulsar 二进制协议生产者/消费者BROKER_HTTP_PORT8080HTTP 管理接口Pulsar Admin说明在org.testcontainers.containers包下还存在一个标有Deprecated的旧版PulsarContainer见 org/testcontainers/containers/PulsarContainer.java其 Javadoc 明确要求改用org.testcontainers.pulsar.PulsarContainer本文一律以新版为准。从源码可见构造器会对镜像名做兼容性校验dockerImageName.assertCompatibleWith(...)仅允许apachepulsar/pulsar与apachepulsar/pulsar-all两个仓库同时完成端口暴露与等待策略的初始化。测试 CompatibleApachePulsarImageTest.java 以参数化方式分别验证了这两种镜像在普通收发与事务场景下的可用性。快速上手创建容器并获取连接地址在测试中以 try-with-resources 创建并启动容器即可容器会在测试结束后自动清理// constructorWithVersion { PulsarContainer pulsar new PulsarContainer(apachepulsar/pulsar:3.0.0); // } pulsar.start();构造器重载既支持直接传镜像名字符串也支持传DockerImageName对象例如DockerImageName.parse(apachepulsar/pulsar:3.0.0)。容器启动后即可获取两条关键连接地址// coordinates { final String pulsarBrokerUrl pulsar.getPulsarBrokerUrl(); final String httpServiceUrl pulsar.getHttpServiceUrl(); // }两个 getter 的实现在 PulsarContainer.java 中一目了然public String getPulsarBrokerUrl() { return String.format(pulsar://%s:%s, getHost(), getMappedPort(BROKER_PORT)); } public String getHttpServiceUrl() { return String.format(http://%s:%s, getHost(), getMappedPort(BROKER_HTTP_PORT)); }getHost()与getMappedPort()返回的是宿主机视角的地址因此无论在本地还是 CI 环境中测试代码都不必关心容器的实际 IP 与随机映射端口。pulsarBrokerUrl用于构造PulsarClienthttpServiceUrl用于构造PulsarAdmin。完整的收发验证示例测试基类 AbstractPulsar.java 给出了一个可复制的端到端验证模板——先启动消费者订阅test_topic再发送消息并异步接收断言PulsarClient client PulsarClient.builder().serviceUrl(pulsarBrokerUrl).build(); Consumerbyte[] consumer client .newConsumer() .topic(test_topic) .subscriptionName(test-subs) .subscribe(); Producerbyte[] producer client.newProducer().topic(test_topic).create(); producer.send(test containers.getBytes()); CompletableFutureMessagebyte[] future consumer.receiveAsync(); Messagebyte[] message future.get(5, TimeUnit.SECONDS); assertThat(new String(message.getData())).isEqualTo(test containers);通过PULSAR_PREFIX_注入任意 Pulsar 配置Pulsar 原生支持环境变量驱动配置standalone 镜像在启动时会执行/pulsar/bin/apply-config-from-env.py /pulsar/conf/standalone.conf把以PULSAR_PREFIX_为前缀的环境变量翻译为standalone.conf中的同名配置项。Testcontainers 直接复用这一机制所以你可以用容器 API 的withEnv(...)设置任意 Pulsar 配置变量例如开启 broker 级消息去重// constructorWithEnv { PulsarContainer pulsar new PulsarContainer(PULSAR_IMAGE) .withEnv(PULSAR_PREFIX_brokerDeduplicationEnabled, true); // }即PULSAR_PREFIX_brokerDeduplicationEnabledtrue会被映射为standalone.conf里的brokerDeduplicationEnabledtrue。自定义集群名是这一机制的典型用例。setupCommandAndEnv()中会读取环境变量getEnvMap().getOrDefault(PULSAR_PREFIX_clusterName, standalone)默认集群名为standalone通过pulsar.withEnv(PULSAR_PREFIX_clusterName, tc-cluster);即可将集群改名为tc-cluster测试 customClusterName 专门验证了该场景。注意集群名同时被用于就绪检查见下文就绪等待策略因此若修改集群名等待逻辑会自动跟随该环境变量无需额外配置。启用 Pulsar Functions WorkerPulsar IOPulsar IO 框架连接器 Source/Sink依赖 Functions Worker。模块默认不启用Functions WorkershouldNotEnableFunctionsWorkerByDefault测试证明默认启动的容器调用pulsarAdmin.functions().getFunctions(...)会抛出PulsarAdminException。如需测试 Pulsar IO链式调用withFunctionsWorker()// constructorWithFunctionsWorker { PulsarContainer pulsar new PulsarContainer(DockerImageName.parse(apachepulsar/pulsar:3.0.0)) .withFunctionsWorker(); // }启用后容器启动命令会去掉--no-functions-worker -nss参数源码见 PulsarContainer.java并额外等待日志中出现.*Function worker service started.*才认为就绪。shouldWaitForFunctionsWorkerStarted测试验证了启用后getFunctions(public, default)可以正常返回大小为 0的空结果。启用 Pulsar TransactionsPulsar 事务能力在 standalone 模式下默认关闭需要通过withTransactions()显式开启// constructorWithTransactions { PulsarContainer pulsar new PulsarContainer(PULSAR_IMAGE).withTransactions(); // }其底层实现见 PulsarContainer.java会做两件事注入环境变量PULSAR_PREFIX_transactionCoordinatorEnabledtrue开启事务协调器增加一条就绪检查轮询GET /admin/v2/persistent/pulsar/system/transaction_coordinator_assign/partitions直到返回 HTTP 200确保事务协调器分配系统 Topic 已经创建。测试 testTransactions 先通过PulsarAdmin断言系统 Topicpersistent://pulsar/system/transaction_coordinator_assign已存在再执行真实的客户端事务流程构建PulsarClient时调用.enableTransaction(true)创建事务、发送消息、提交最后消费并断言消息内容完整流程见 AbstractPulsar.javafinal Transaction transaction client.newTransaction().build().get(); producer.newMessage(transaction).value(first).send(); transaction.commit(); MessageString message consumer.receive(); assertThat(message.getValue()).isEqualTo(first);withTransactions()与withFunctionsWorker()可以同时调用testTransactionsAndFunctionsWorker测试即验证了两者共存时的完整功能。就绪等待策略容器何时算真正可用这是本模块最容易踩坑也最值得关注的部分。启动命令与等待策略都封装在 PulsarContainer.java 的setupCommandAndEnv()中String standaloneBaseCommand /pulsar/bin/apply-config-from-env.py /pulsar/conf/standalone.conf bin/pulsar standalone; if (!functionsWorkerEnabled) { standaloneBaseCommand --no-functions-worker -nss; } withCommand(/bin/bash, -c, standaloneBaseCommand);容器以/bin/bash -c依次执行应用环境变量配置与启动 standalone 服务默认追加--no-functions-worker -nss以跳过 Functions Worker 与相关系统服务从而加快启动速度只有启用withFunctionsWorker()时才去掉该参数。就绪判定采用WaitAllStrategy按开关叠加多条策略触发条件等待策略端点 / 断言总是启用HTTP 轮询GET /admin/v2/clusters响应体须等于[clusterName]默认[standalone]证明集群已初始化完成withTransactions()HTTP 轮询GET /admin/v2/persistent/pulsar/system/transaction_coordinator_assign/partitions返回 200withFunctionsWorker()日志匹配日志出现.*Function worker service started.*其中集群检查对应测试testClusterFullyInitialized它断言pulsarAdmin.clusters().getClusters()恰好包含一个名为standalone的集群。此外testStartupTimeoutIsHonored表明启动超时可通过标准的withStartupTimeout(Duration)控制。将模块加入项目依赖在pom.xml/build.gradle中添加如下依赖版本号{{latest_version}}以当前最新发布版本为准 Gradlegroovy testImplementation org.testcontainers:testcontainers-pulsar:{{latest_version}} Mavenxml dependency groupIdorg.testcontainers/groupId artifactIdtestcontainers-pulsar/artifactId version{{latest_version}}/version scopetest/scope /dependency从模块构建文件 modules/pulsar/build.gradle 可以看到模块以api方式依赖核心库:testcontainers因此仅引入本模块即可获得GenericContainer等全部基础能力。测试侧还需要 Pulsar 客户端依赖仓库本身通过 BOM 统一管理版本testImplementation platform(org.apache.pulsar:pulsar-bom:4.2.0) testImplementation org.apache.pulsar:pulsar-client testImplementation org.apache.pulsar:pulsar-client-admin在你的项目中如需编写收发与事务相关测试同样需要引入pulsar-client构造PulsarClient、Producer、Consumer与pulsar-client-admin构造PulsarAdmin这两个 artifact。典型测试场景小结综合上述能力Pulsar 模块可以在 JUnit 测试中覆盖以下场景全部无需外部服务消息收发默认容器即可验证 Producer/Consumer 的发布订阅语义管理 API通过httpServiceUrl构造PulsarAdmin验证 Topic、集群、函数等管理操作Pulsar IO / 连接器withFunctionsWorker()启用 Functions Worker 后测试 Source/Sink 流程事务withTransactions()开启事务协调器验证事务性生产与提交配置实验借助PULSAR_PREFIX_环境变量机制无需改动代码即可实验任意 broker 配置例如brokerDeduplicationEnabled、自定义clusterName等。需要注意的是容器使用 standalone 单节点模式适合功能与集成测试若需要验证多节点集群行为则应结合 Testcontainers 的能力另行构建多容器拓扑而非依赖本模块的单容器封装。此外运行本模块测试的前置条件是本地具备可用的 Docker 环境这也是整个 Testcontainers 项目的通用前提。【免费下载链接】testcontainers-javaTestcontainers is a Java library that supports JUnit tests, providing lightweight, throwaway instances of common databases, Selenium web browsers, or anything else that can run in a Docker container.项目地址: https://gitcode.com/GitHub_Trending/te/testcontainers-java创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考