如何在单个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
相关产品推荐
相关产品推荐

