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

如何在单个Ballerina项目中为一个主题创建两个订阅者?

在单个Ballerina项目中为主题创建两个订阅者的解决方案

问题描述

需求:在单个Ballerina项目中为一个主题创建两个独立订阅者(subscriber1、subscriber2)。

已实现的思路:

  • 创建init()函数,为每个订阅者初始化ASBServiceReceiverConfig实例receiver1和receiver2
  • 在public main函数中接收每个接收器的消息,存储到asb:Message类型变量msg1、msg2
  • 为两个订阅者分别编写function1、function2,函数内包含调用外部REST API的逻辑

当前困境:
现有循环逻辑仅能处理单个订阅者,无法同时运行两个订阅者的消息处理流程;尝试使用workers但未成功触发API调用,且读取一条消息后就停止,不知道如何让处理线程持续运行。

现有单订阅者循环代码:

while true {
do {
    if (msg1 is ()) {
        // 0 messages to read in the topic
        continue;
    } 
    //Call the function which has logic to call an API
    function1(receiver1, msg1);    
} on fail error e{
            // log error here
    }
}

解决方案

方法1:使用Ballerina Workers实现并行持续处理

Ballerina Workers支持并行执行逻辑,之前的问题大概率是未在worker内部编写持续循环的逻辑。正确写法是给每个订阅者分配独立worker,在worker内部实现持续读取并处理消息的循环:

import ballerina/asb;
import ballerina/io;
import ballerina/log;

public function main() returns error? {
    // 假设已在init()中完成接收器初始化,直接引用实例
    asb:ServiceReceiverConfig receiver1 = ...;
    asb:ServiceReceiverConfig receiver2 = ...;

    // 启动worker处理subscriber1
    worker worker1 {
        while true {
            do {
                asb:Message? msg1 = asb->receive(receiver1);
                if (msg1 is ()) {
                    // 无消息时休眠1秒,避免空循环占用CPU
                    io:sleep(1000);
                    continue;
                }
                // 调用订阅者1的业务逻辑函数
                check function1(receiver1, msg1);
            } on fail error e {
                log:printError("Subscriber1处理失败", err = e);
                // 出错后休眠2秒再重试
                io:sleep(2000);
            }
        }
    }

    // 启动worker处理subscriber2
    worker worker2 {
        while true {
            do {
                asb:Message? msg2 = asb->receive(receiver2);
                if (msg2 is ()) {
                    io:sleep(1000);
                    continue;
                }
                check function2(receiver2, msg2);
            } on fail error e {
                log:printError("Subscriber2处理失败", err = e);
                io:sleep(2000);
            }
        }
    }

    // 阻塞main函数,确保workers持续运行
    wait worker1;
    wait worker2;
}

// Subscriber1的业务逻辑:调用外部REST API
function function1(asb:ServiceReceiverConfig receiver, asb:Message msg) returns error? {
    import ballerina/http;
    http:Client apiClient = check new("https://api.example.com/subscriber1");
    http:Response res = check apiClient->post("/process", msg.body);
    // 确认消息已处理(根据ASB配置决定是否需要)
    check asb->complete(receiver, msg);
}

// Subscriber2的业务逻辑:调用另一个REST API
function function2(asb:ServiceReceiverConfig receiver, asb:Message msg) returns error? {
    import ballerina/http;
    http:Client apiClient = check new("https://api.example.com/subscriber2");
    http:Response res = check apiClient->put("/handle", msg.body);
    check asb->complete(receiver, msg);
}

关键要点:

  • 每个worker内部包含独立的while true循环,实现持续监听消息
  • 无消息时添加io:sleep()避免空循环消耗过多CPU资源
  • 异常处理中添加休眠时间,防止频繁报错
  • 最后用wait语句阻塞main函数,防止进程退出导致workers终止

方法2:使用start启动异步任务

如果workers写法过于繁琐,可使用start关键字(Ballerina的spawn语法糖)启动异步任务,每个任务对应一个订阅者的持续处理逻辑:

import ballerina/asb;
import ballerina/io;
import ballerina/log;
import ballerina/http;

public function main() returns error? {
    asb:ServiceReceiverConfig receiver1 = ...;
    asb:ServiceReceiverConfig receiver2 = ...;

    // 启动异步任务处理subscriber1
    _ = start runSubscriber(receiver1, handleSubscriber1);
    // 启动异步任务处理subscriber2
    _ = start runSubscriber(receiver2, handleSubscriber2);

    // 让main函数持续运行,直到手动终止进程
    io:waitForExit();
}

// 通用订阅者处理函数,封装循环逻辑
function runSubscriber(asb:ServiceReceiverConfig receiver, function(asb:ServiceReceiverConfig, asb:Message) returns error? handler) returns error? {
    while true {
        do {
            asb:Message? msg = asb->receive(receiver);
            if (msg is ()) {
                io:sleep(1000);
                continue;
            }
            check handler(receiver, msg);
        } on fail error e {
            log:printError("订阅者处理失败", err = e);
            io:sleep(2000);
        }
    }
}

// Subscriber1的API调用逻辑
function handleSubscriber1(asb:ServiceReceiverConfig receiver, asb:Message msg) returns error? {
    http:Client client = check new("https://api.example.com/sub1");
    http:Response res = check client->post("/process", msg.body);
    check asb->complete(receiver, msg);
}

// Subscriber2的API调用逻辑
function handleSubscriber2(asb:ServiceReceiverConfig receiver, asb:Message msg) returns error? {
    http:Client client = check new("https://api.example.com/sub2");
    http:Response res = check client->put("/handle", msg.body);
    check asb->complete(receiver, msg);
}

关键要点:

  • 用start启动异步任务,两个订阅者的处理逻辑并行执行
  • 提取通用的runSubscriber函数,减少重复代码
  • 使用io:waitForExit()让main函数保持运行,避免进程提前终止
  • 处理完消息后调用asb->complete()确认消息(根据ASB的配置要求)

内容的提问来源于stack exchange,提问作者Maryam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 07:35:42