You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在嵌入式独立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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:27:04