如何在嵌入式独立Apache Pulsar环境中注册并测试Functions
在嵌入式Pulsar环境中注册自定义Function并完成集成测试
我来帮你搞定在嵌入式Pulsar环境里注册自定义Function的集成测试问题!要实现这个需求,我们需要先开启Pulsar的Functions Worker组件,再通过Admin API注册函数,最后验证函数的处理逻辑是否正常。下面结合你的Spock测试代码,一步步来实现:
步骤1:修改嵌入式Pulsar配置,启用Functions Worker
默认的PulsarStandalone不会自动启动Functions Worker,这是运行Pulsar Functions的核心组件,所以我们需要先修改配置开启它:
def setupSpec() { def conf = new ServiceConfiguration( clusterName: "test-cluster", zookeeperServers: "localhost:2184", functionsWorkerEnabled: true, // 开启Functions Worker服务 functionsWorkerServiceUrl: "http://localhost:8080", // Functions Worker的API访问地址 brokerServiceUrl: "pulsar://localhost:6650" // 明确指定Broker地址 ) log.info("Starting embedded Pulsar with Functions Worker enabled") standalone = PulsarStandaloneBuilder.instance() .withConfig(conf) .withNoStreamStorage(true) .build() standalone.start() // 直接从standalone实例获取PulsarService,无需单独创建 pulsarService = standalone.pulsarService }
步骤2:编写自定义Pulsar Function类
创建一个实现org.apache.pulsar.functions.api.Function的处理类,比如做一个简单的字符串转大写逻辑:
package kic.data.stream.pulsar.functions import org.apache.pulsar.functions.api.Function class StringToUpperCaseFunction implements Function<String, String> { @Override String process(String input, Context context) throws Exception { return input != null ? input.toUpperCase() : null } }
步骤3:初始化Pulsar Admin并注册函数
我们需要通过PulsarAdmin来管理函数的生命周期(注册、删除),在测试类中添加相关逻辑:
首先在测试类顶部添加静态变量:
static final String INPUT_TOPIC = "persistent://public/default/input-topic" static final String OUTPUT_TOPIC = "persistent://public/default/output-topic" static final String FUNCTION_NAME = "string-upper-case-function" static PulsarAdmin pulsarAdmin
然后在setupSpec中初始化Admin并注册函数:
def setupSpec() { // ... 前面的PulsarStandalone启动代码 ... // 初始化PulsarAdmin客户端 pulsarAdmin = PulsarAdmin.builder() .serviceHttpUrl("http://localhost:8080") // 和Functions Worker地址保持一致 .build() // 构建函数配置 def functionConfig = new FunctionConfig() functionConfig.setName(FUNCTION_NAME) functionConfig.setTenant("public") functionConfig.setNamespace("default") functionConfig.setInputTopics(Collections.singleton(INPUT_TOPIC)) functionConfig.setOutputTopic(OUTPUT_TOPIC) functionConfig.setClassName(StringToUpperCaseFunction.class.getName()) functionConfig.setRuntime(Runtime.JAVA) functionConfig.setParallelism(1) // 避免重复注册:如果函数已存在则先删除 try { pulsarAdmin.functions().deleteFunction("public", "default", FUNCTION_NAME) } catch (NotFoundException e) { // 函数不存在,忽略异常 } // 注册自定义函数 pulsarAdmin.functions().createFunction(functionConfig) }
步骤4:修改测试逻辑,验证函数处理结果
调整测试用例,让生产者发送消息到输入主题,消费者从输出主题接收处理后的消息,验证逻辑是否符合预期:
def "test pulsar function message processing"() { given: PulsarClient client = PulsarClient.builder() .serviceUrl(pulsarService.brokerServiceUrl) .build() Producer<String> producer = client.newProducer(Schema.STRING) .topic(INPUT_TOPIC) .enableBatching(false) .create() Consumer<String> consumer = client.newConsumer(Schema.STRING) .topic(OUTPUT_TOPIC) .subscriptionName("function-test-subs") .ackTimeout(10, TimeUnit.SECONDS) .subscriptionType(SubscriptionType.Exclusive) .subscribe() when: // 发送100条测试消息 def testMessages = (1..NUM_OF_MESSAGES).collect { "Hello_$it" } testMessages.each { producer.send(it) } // 接收处理后的消息 def receivedMessages = [] for (int i = 1; i <= NUM_OF_MESSAGES; ++i) { Message<String> message = consumer.receive(5, TimeUnit.SECONDS) receivedMessages.add(message.value()) consumer.acknowledge(message) } then: // 验证所有消息都被转成大写 receivedMessages == testMessages.collect { it.toUpperCase() } cleanup: producer.close() consumer.close() client.close() }
步骤5:添加测试后清理逻辑
在cleanupSpec中删除注册的函数并关闭Admin客户端,避免资源残留:
def cleanupSpec() { // 删除测试函数 try { pulsarAdmin.functions().deleteFunction("public", "default", FUNCTION_NAME) } catch (Exception e) { log.warning("Failed to delete test function: ${e.message}") } pulsarAdmin.close() standalone.close() }
注意事项
- 确保项目依赖中包含
pulsar-functions-api和pulsar-client-admin包,否则会出现类找不到的问题 - 如果测试中出现消息接收超时,可以适当延长
receive方法的超时时间,给函数足够的启动和处理时间 - 自定义Function类要确保在测试的classpath范围内
内容的提问来源于stack exchange,提问作者KIC
相关产品推荐
相关产品推荐

