pbai 发表于 7 天前

PBIDEA:用 uo_redis 与 uo_kafka 做消息队列与发布订阅

PBIDEA:用 uo_redis 与 uo_kafka 做消息队列与发布订阅


阅读说明
1. 适用版本:PBIDEA(本机运行库 PbIdea.dll;客户端对象位于 websuite.pbl / pbjson.pbl,对应 PBIDEA 1.x,建议 PB 10 及以上宿主环境)
2. 支持数据库:本文不涉及关系型数据库;依赖中间件 Redis(建议 5.0+)与 Apache Kafka(建议 2.x+)作为外部服务
3. 操作系统与环境要求:Windows 7+,已安装 PBIDEA 运行库(PbIdea.dll),且运行机器能访问 Redis / Kafka 服务端口(默认 6379 / 9092)
4. 难度系数:★★★★☆(需理解消息队列与发布订阅模型,并配置外部中间件)
5. 其它阅读说明:示例默认 Redis 无密码、Kafka 单 broker(localhost:9092);生产环境请配置密码/集群并处理好异常


一、为什么要上消息队列

很多 PB 老系统里,模块之间是"直接调用":下单成功后立刻发短信、立刻写统计、立刻通知库存。一旦其中一环慢或挂了,整个下单就被拖死。消息队列把"做一件事"和"这件事引发的一堆后续动作"解耦:


[*]异步:主流程只把消息投出去就返回,后续动作由消费者慢慢处理。
[*]削峰:突发流量先堆在队列里,消费者按自己的节奏消费,不会压垮下游。
[*]解耦:生产者不认识消费者,新增一个订阅方不用改生产方代码。


本文用 PBIDEA 自带的两个客户端对象把这套能力接进 PB:uo_redis(同步 + 异步)和 uo_kafka_producer / uo_kafka_consumer。

二、两套方案怎么选


维度Redis(uo_redis)Kafka(uo_kafka_producer / uo_kafka_consumer)
定位轻量、低延迟,单机/主从即可分布式、高吞吐、持久化、可水平扩展
队列模型基于 List(LPUSH/BRPOP)做点对点;Pub/Sub 做广播基于 Topic + 分区,天然支持多消费者组、可重放
适用规模中小流量、内部系统解耦大规模事件流、日志流、跨系统总线
部署成本一个 redis-server 即可需 Zookeeper/Controller + broker 集群


一句话:轻量内部解耦用 Redis,规模化事件总线用 Kafka。下面分别给出可运行示例。

三、uo_redis 同步客户端对象全貌

uo_redis 是 nonvisualobject,构造时自动 redisCreate()、析构时自动 Close(),所有方法最终走 PbIdea.dll。常用 API:


[*]boolean Open(string ip, long port, string auth):连接;auth 传空串 '' 表示无密码。
[*]Close():断开。
[*]uo_redisreply Set(string key, string/long/double value):写入键值对(自动转 UTF-8)。
[*]uo_redisreply Get(string key):取键,返回 uo_redisreply 对象。
[*]boolean Get(string key, ref string/long/double value):直接把值取出到变量,省去解析 reply。
[*]uo_redisreply ExecCommand(string fmt, ...):执行任意 Redis 命令,支持 %s/%d 占位与变长参数——队列/LIST 操作全靠它。
[*]BeginPipeline() / EndPipeline(ref uo_redisreply replys[]) / EndPipeline(ref uo_json replys) / AppendCommand(string fmt, ...):流水线批量,减少往返。


返回结果统一包装成 uo_redisreply,它有两个核心方法:


[*]boolean Get(ref string/blob/long/uo_redisreply values[]):把值取出来;数组类型用 ref uo_redisreply values[] 取子元素。
[*]int GetType():返回类型常量,如 REPLY_STRING=1、REPLY_ARRAY=2、REPLY_INTEGER=3、REPLY_NIL=4、REPLY_STATUS=5、REPLY_ERROR=6、REPLY_DOUBLE=7 等。判类型就靠这些常量,避免写魔法数字。


四、示例 1:用 Redis List 做点对点队列(生产者)

思路:生产者用 LPUSH 把消息压入列表头,消费者用 BRPOP 从列表尾阻塞弹出——先进先出,一条消息只会被一个消费者拿走。


前置:本机已启动 Redis 服务(默认 127.0.0.1:6379)且无密码;在窗口按钮的 clicked 事件贴入下方代码即可运行。
步骤:1) 拖一个按钮 cb_push;2) 在 clicked 事件粘贴代码;3) 运行并点击,MessageBox 显示"入队成功"。



// 示例输入:Redis 服务地址
string ls_ip
ls_ip = '127.0.0.1'
// 示例输入:Redis 服务端口
long ll_port
ll_port = 6379
// 示例输入:认证密码(无密码填空串)
string ls_auth
ls_auth = ''
// 示例输入:队列名
string ls_queue
ls_queue = 'pb_task_queue'
// 示例输入:要投递的消息
string ls_msg
ls_msg = 'order_1001_paid'

uo_redis hc
hc = create uo_redis
boolean lb_ok
lb_ok = hc.Open(ls_ip, ll_port, ls_auth)
if not lb_ok then
    MessageBox('连接失败', '请检查 Redis 服务是否启动')
    destroy hc
    return
end if

// 入队:LPUSH 把消息压入列表头部,%s 依次对应 ls_queue、ls_msg
uo_redisreply rp
rp = hc.ExecCommand('LPUSH %s %s', ls_queue, ls_msg)
int li_type
li_type = rp.GetType()
if li_type = rp.REPLY_STATUS then
    // LPUSH 成功返回 +OK 状态
    MessageBox('入队成功', '队列=' + ls_queue + ' 消息=' + ls_msg)
else
    MessageBox('入队异常', 'type=' + string(li_type))
end if

hc.Close()
destroy hc


五、示例 2:消费者出队并解析 reply(数组取值)

BRPOP 返回的是一个数组 [队列名, 消息内容],所以要用 GetType() 判 REPLY_ARRAY,再 Get(ref reply_array[]) 取出子元素。注意 PB 数组是 1 基,reply_array 才是真正的消息体。


前置:同上,Redis 服务已启动;本例从队列取一条消息并打印。



// 示例输入:Redis 地址
string ls_ip
ls_ip = '127.0.0.1'
// 示例输入:端口
long ll_port
ll_port = 6379
// 示例输入:密码(无则空串)
string ls_auth
ls_auth = ''
// 示例输入:队列名
string ls_queue
ls_queue = 'pb_task_queue'
// 示例输入:阻塞超时(秒),0 表示一直阻塞直到取到消息
long ll_timeout
ll_timeout = 5

uo_redis hc
hc = create uo_redis
boolean lb_ok
lb_ok = hc.Open(ls_ip, ll_port, ls_auth)
if not lb_ok then
    MessageBox('连接失败', '请检查 Redis 服务')
    destroy hc
    return
end if

// 出队:BRPOP 阻塞弹出列表尾部元素;%d 对应 ll_timeout
uo_redisreply rp
rp = hc.ExecCommand('BRPOP %s %d', ls_queue, ll_timeout)
int li_type
li_type = rp.GetType()
if li_type = rp.REPLY_ARRAY then
    uo_redisreply reply_array[]
    rp.Get(ref reply_array)          // 取出数组:=队列名, =消息体
    string ls_val
    reply_array.Get(ref ls_val)   // 第 2 个元素是真正的消息内容
    MessageBox('收到消息', ls_val)
elseif li_type = rp.REPLY_NIL then
    // 超时仍未取到任何消息
    MessageBox('队列空', '等待超时,未取到消息')
else
    MessageBox('取消息异常', 'type=' + string(li_type))
end if

hc.Close()
destroy hc


六、简单取值的捷径

如果只是简单的"存一个字符串、取一个字符串",不必走 ExecCommand 再解析数组,uo_redis 提供了直取重载:


// 示例输入:键与待存的值
string ls_key
ls_key = 'pb_config:env'
string ls_put
ls_put = 'prod'
// 取出用 ref 变量承接
string ls_get
ls_get = ''

uo_redis hc
hc = create uo_redis
hc.Open('127.0.0.1', 6379, '')
hc.Set(ls_key, ls_put)               // 自动转 UTF-8 存储
boolean lb_got
lb_got = hc.Get(ls_key, ref ls_get)// 直接取到 ls_get,返回是否取到
MessageBox('读取结果', ls_key + ' = ' + ls_get + ' (got=' + string(lb_got) + ')')
hc.Close()
destroy hc


七、Redis 发布订阅(异步 uo_redisasync)

上面的 List 是"点对点"(一条消息一个消费者)。如果要"一份消息广播给多个订阅方",用 Redis 的 Pub/Sub。同步 uo_redis 没有订阅方法,要用异步版 uo_redisasync:连接后 Subscribe(频道, true) 开启订阅,收到推送时触发 onreply(unsignedlong id, integer replycount, readonly string replys[]) 事件。


// 前置:需要长连接并接收推送的场景;声明实例变量 ha 并编写 onreply 事件处理。
// 示例输入:Redis 地址
string ls_ip
ls_ip = '127.0.0.1'
// 示例输入:端口
long ll_port
ll_port = 6379
// 示例输入:频道名
string ls_channel
ls_channel = 'pb_notify'

uo_redisasync ha
ha = create uo_redisasync
boolean lb_ok
lb_ok = ha.Open(ls_ip, ll_port)
if not lb_ok then
    MessageBox('连接失败', '请检查 Redis 服务')
    destroy ha
    return
end if
ulong lul_sub
lul_sub = ha.Subscribe(ls_channel, true)   // switch=true 开启订阅,返回订阅 id
MessageBox('订阅已提交', 'channel=' + ls_channel + ' id=' + string(lul_sub))
// 收到推送时,ha 的 onreply 事件被触发,在其中处理 replys[] 即可


onreply 事件里按 replys[] 取出每条推送内容;配合 onconnect/ondisconnect 事件可以做断线提示与重连。

八、uo_kafka_producer 全貌与示例

Kafka 适合"高吞吐、可重放、多消费者组"的事件流。uo_kafka_producer 构造时自动 ProducerCreate()、析构自动 ProducerDestroy()。常用 API:


[*]ProducerConnect(string brokers, long timeoutms):连接 broker(多个用逗号分隔,如 'host1:9092,host2:9092')。
[*]ProducerConfigSet(string key, string value) / ProducerConfigGet(string key):设置/读取配置项。
[*]ProducerReadConfig(string section, string configFile):从配置文件按段读取配置(可选)。
[*]int ProducerSend(string topic, string text) / (topic, blob data):发送;返回 int,0 表示成功,非 0 用 ProducerGetError() 取错误。
[*]string ProducerGetError():取最近一次错误详情。
[*]ProducerDisconnect():断开。



// 前置:已部署 Kafka(单 broker,localhost:9092);窗口按钮 clicked 事件贴入下方代码运行。
// 示例输入:Kafka broker 地址(多个用逗号分隔)
string ls_brokers
ls_brokers = 'localhost:9092'
// 示例输入:连接超时(毫秒)
long ll_timeout
ll_timeout = 5000
// 示例输入:目标主题
string ls_topic
ls_topic = 'pb_events'
// 示例输入:要发送的消息
string ls_msg
ls_msg = 'user_login|1001'

uo_kafka_producer kp
kp = create uo_kafka_producer
kp.ProducerConnect(ls_brokers, ll_timeout)
int li_rtn
li_rtn = kp.ProducerSend(ls_topic, ls_msg)   // 也可传 blob 发二进制
if li_rtn = 0 then
    MessageBox('发送成功', 'topic=' + ls_topic)
else
    MessageBox('发送失败', kp.ProducerGetError())
end if
kp.ProducerDisconnect()
destroy kp


九、uo_kafka_consumer 全貌与示例

uo_kafka_consumer 通过 ConsumerPoll 驱动消息分发,消息到达时触发 ue_message(string status, blob data) 事件。常用 API:


[*]ConsumerConnect(string brokers, long timeoutms) / ConsumerDisconnect()
[*]ConsumerSubscribe(string topic):订阅主题,返回 int(0 成功)。
[*]ConsumerUnsubscribe():取消订阅。
[*]boolean ConsumerPoll(long timeout):拉取一次,有消息时触发 ue_message;返回值表示是否还有更多待处理消息。需在 timer 或后台线程里循环调用。
[*]ConsumerConfigSet/Get、ConsumerGetError()



前置:在窗口(或 NVO)声明实例变量 iuc,并编写 iuc 的 ue_message 事件;启动后持续 Poll 即可收到消息。



// 1) 窗口/对象声明实例变量(实例作用域)
uo_kafka_consumer iuc

// 2) iuc 的 ue_message 事件(收到消息时触发):status 为状态串,data 为消息体(blob)
string ls_status
ls_status = status
// data 是 blob,按业务解码(例如转字符串后处理)
MessageBox('收到消息', ls_status)

// 3) 启动消费(按钮 clicked 事件)
// 示例输入:broker 地址
string ls_brokers
ls_brokers = 'localhost:9092'
// 示例输入:连接超时(毫秒)
long ll_timeout
ll_timeout = 5000
// 示例输入:订阅主题
string ls_topic
ls_topic = 'pb_events'

iuc = create uo_kafka_consumer
iuc.ConsumerConnect(ls_brokers, ll_timeout)
int li_sub
li_sub = iuc.ConsumerSubscribe(ls_topic)
if li_sub = 0 then
    MessageBox('订阅成功', 'topic=' + ls_topic)
else
    MessageBox('订阅失败', iuc.ConsumerGetError())
end if
// 接着在窗口 timer 事件里循环调用 iuc.ConsumerPoll(1000) 驱动消息分发


窗口 timer 事件(先 timer(1) 让窗口每秒触发一次):


// 窗口 timer 事件:周期性拉取,消息经 ue_message 回调
boolean lb_more
lb_more = iuc.ConsumerPoll(1000)
// lb_more=true 表示还有积压消息,可继续 Poll;false 表示本次已空


十、边界情况与常见坑


[*]Redis 无密码:Open 的 auth 必须传空串 '',传 null 或省略会连不上。
[*]BRPOP 是阻塞的:timeout=0 会一直阻塞当前调用线程。生产环境务必设合理超时(如 5~30 秒),或把消费放到 uo_thread 后台线程,避免卡 UI。
[*]reply 数组下标:BRPOP 返回 [队列名, 值],取消息体用 reply_array(PB 数组 1 基),别写反。
[*]reply 对象释放:Get/ExecCommand 返回的 uo_redisreply 由析构自动 replyDestroy(),函数返回后随引用失效而回收;不要在长循环里长期 hold 大量 reply 不释放。
[*]Kafka 发送必须判错:ProducerSend 返回非 0 时一定调 ProducerGetError() 看原因(常见:broker 不可达、topic 不存在、序列化失败)。
[*]ConsumerPoll 必须循环:只在 clicked 里调一次只能拿到瞬时消息;消息是靠 ue_message 事件异步回调的,必须由 timer/线程持续 Poll 才能不断接收。
[*]断线处理:网络抖动会被吞掉,建议监听 ProducerGetError / ConsumerGetError,在 UI 上提示并支持重连;uo_redisasync 自带 onconnect/ondisconnect 事件便于感知。
[*]ExecCommand 占位符顺序:fmt 里的 %s/%d 要与后面变长参数一一对应,类型也要匹配(数字用 %d、字符串用 %s),否则命令会被发错。


十一、与相近方案对比 & 进阶扩展


[*]Redis List vs Pub/Sub:List 是点对点、可堆积(消费者离线消息不丢);Pub/Sub 是广播、不堆积(订阅方离线期间的消息收不到)。按"要不要留存"选。
[*]Redis vs Kafka:小规模、低延迟、单点够用选 Redis;要持久化、可重放、多消费者组、跨机房选 Kafka。
[*]Pipeline 批量:高频写场景用 BeginPipeline() + AppendCommand(...) + EndPipeline(ref replys[]),一次网络往返执行多条命令,吞吐显著提升。
[*]异步订阅:需要服务端主动推送用 uo_redisasync,连接后 Subscribe 并用 onreply 处理,比同步轮询更省资源。
[*]与界面/数据联动:消费者在 ue_message / onreply 里拿到消息后,调用 dw_1.Retrieve() 或 dw_1.InsertRow() 刷新界面;耗时处理建议放进 uo_thread(见 Day 10)后台执行,避免阻塞主线程。


把这套接好,PB 老系统就能以很低的成本获得"异步 + 削峰 + 解耦"的能力,而不必重写架构。
页: [1]
查看完整版本: PBIDEA:用 uo_redis 与 uo_kafka 做消息队列与发布订阅

免责声明:
本站所发布的一切破解补丁、注册机和注册信息及软件的解密分析文章仅限用于学习和研究目的;不得将上述内容用于商业或者非法用途,否则,一切后果请用户自负。本站信息来自网络,版权争议与本站无关。您必须在下载后的24个小时之内,从您的电脑中彻底删除上述内容。如果您喜欢该程序,请支持正版软件,购买注册,得到更好的正版服务。如有侵权请邮件与我们联系处理。

Mail To:Admin@SybaseBbs.com