祝愿大家身体健康!

 站点注册  找回密码
 站点注册

QQ登录

只需一步,快速开始

查看: 82|回复: 0

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

[复制链接]

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

[复制链接]
pbai

主题

0

回帖

1348

积分

PBAI

积分
1348
贡献
在线时间
小时
昨天 07:19 | 显示全部楼层 |阅读模式

马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。

您需要 登录 才可以下载或查看,没有账号?站点注册

×
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=1REPLY_ARRAY=2REPLY_INTEGER=3REPLY_NIL=4REPLY_STATUS=5REPLY_ERROR=6REPLY_DOUBLE=7 等。判类型就靠这些常量,避免写魔法数字。


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

思路:生产者用 LPUSH 把消息压入列表头,消费者用 BRPOP 从列表尾阻塞弹出——先进先出,一条消息只会被一个消费者拿走。
前置:本机已启动 Redis 服务(默认 127.0.0.1:6379)且无密码;在窗口按钮的 clicked 事件贴入下方代码即可运行。
步骤:1) 拖一个按钮 cb_push;2) 在 clicked 事件粘贴代码;3) 运行并点击,MessageBox 显示"入队成功"。
  1. // 示例输入:Redis 服务地址
  2. string ls_ip
  3. ls_ip = '127.0.0.1'
  4. // 示例输入:Redis 服务端口
  5. long ll_port
  6. ll_port = 6379
  7. // 示例输入:认证密码(无密码填空串)
  8. string ls_auth
  9. ls_auth = ''
  10. // 示例输入:队列名
  11. string ls_queue
  12. ls_queue = 'pb_task_queue'
  13. // 示例输入:要投递的消息
  14. string ls_msg
  15. ls_msg = 'order_1001_paid'
  16. uo_redis hc
  17. hc = create uo_redis
  18. boolean lb_ok
  19. lb_ok = hc.Open(ls_ip, ll_port, ls_auth)
  20. if not lb_ok then
  21.     MessageBox('连接失败', '请检查 Redis 服务是否启动')
  22.     destroy hc
  23.     return
  24. end if
  25. // 入队:LPUSH 把消息压入列表头部,%s 依次对应 ls_queue、ls_msg
  26. uo_redisreply rp
  27. rp = hc.ExecCommand('LPUSH %s %s', ls_queue, ls_msg)
  28. int li_type
  29. li_type = rp.GetType()
  30. if li_type = rp.REPLY_STATUS then
  31.     // LPUSH 成功返回 +OK 状态
  32.     MessageBox('入队成功', '队列=' + ls_queue + ' 消息=' + ls_msg)
  33. else
  34.     MessageBox('入队异常', 'type=' + string(li_type))
  35. end if
  36. hc.Close()
  37. destroy hc
复制代码

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

BRPOP 返回的是一个数组 [队列名, 消息内容],所以要用 GetType()REPLY_ARRAY,再 Get(ref reply_array[]) 取出子元素。注意 PB 数组是 1 基,reply_array[2] 才是真正的消息体。
前置:同上,Redis 服务已启动;本例从队列取一条消息并打印。
  1. // 示例输入:Redis 地址
  2. string ls_ip
  3. ls_ip = '127.0.0.1'
  4. // 示例输入:端口
  5. long ll_port
  6. ll_port = 6379
  7. // 示例输入:密码(无则空串)
  8. string ls_auth
  9. ls_auth = ''
  10. // 示例输入:队列名
  11. string ls_queue
  12. ls_queue = 'pb_task_queue'
  13. // 示例输入:阻塞超时(秒),0 表示一直阻塞直到取到消息
  14. long ll_timeout
  15. ll_timeout = 5
  16. uo_redis hc
  17. hc = create uo_redis
  18. boolean lb_ok
  19. lb_ok = hc.Open(ls_ip, ll_port, ls_auth)
  20. if not lb_ok then
  21.     MessageBox('连接失败', '请检查 Redis 服务')
  22.     destroy hc
  23.     return
  24. end if
  25. // 出队:BRPOP 阻塞弹出列表尾部元素;%d 对应 ll_timeout
  26. uo_redisreply rp
  27. rp = hc.ExecCommand('BRPOP %s %d', ls_queue, ll_timeout)
  28. int li_type
  29. li_type = rp.GetType()
  30. if li_type = rp.REPLY_ARRAY then
  31.     uo_redisreply reply_array[]
  32.     rp.Get(ref reply_array)          // 取出数组:[1]=队列名, [2]=消息体
  33.     string ls_val
  34.     reply_array[2].Get(ref ls_val)   // 第 2 个元素是真正的消息内容
  35.     MessageBox('收到消息', ls_val)
  36. elseif li_type = rp.REPLY_NIL then
  37.     // 超时仍未取到任何消息
  38.     MessageBox('队列空', '等待超时,未取到消息')
  39. else
  40.     MessageBox('取消息异常', 'type=' + string(li_type))
  41. end if
  42. hc.Close()
  43. destroy hc
复制代码

六、简单取值的捷径

如果只是简单的"存一个字符串、取一个字符串",不必走 ExecCommand 再解析数组,uo_redis 提供了直取重载:
  1. // 示例输入:键与待存的值
  2. string ls_key
  3. ls_key = 'pb_config:env'
  4. string ls_put
  5. ls_put = 'prod'
  6. // 取出用 ref 变量承接
  7. string ls_get
  8. ls_get = ''
  9. uo_redis hc
  10. hc = create uo_redis
  11. hc.Open('127.0.0.1', 6379, '')
  12. hc.Set(ls_key, ls_put)               // 自动转 UTF-8 存储
  13. boolean lb_got
  14. lb_got = hc.Get(ls_key, ref ls_get)  // 直接取到 ls_get,返回是否取到
  15. MessageBox('读取结果', ls_key + ' = ' + ls_get + ' (got=' + string(lb_got) + ')')
  16. hc.Close()
  17. destroy hc
复制代码

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

上面的 List 是"点对点"(一条消息一个消费者)。如果要"一份消息广播给多个订阅方",用 Redis 的 Pub/Sub。同步 uo_redis 没有订阅方法,要用异步版 uo_redisasync:连接后 Subscribe(频道, true) 开启订阅,收到推送时触发 onreply(unsignedlong id, integer replycount, readonly string replys[]) 事件。
  1. // 前置:需要长连接并接收推送的场景;声明实例变量 ha 并编写 onreply 事件处理。
  2. // 示例输入:Redis 地址
  3. string ls_ip
  4. ls_ip = '127.0.0.1'
  5. // 示例输入:端口
  6. long ll_port
  7. ll_port = 6379
  8. // 示例输入:频道名
  9. string ls_channel
  10. ls_channel = 'pb_notify'
  11. uo_redisasync ha
  12. ha = create uo_redisasync
  13. boolean lb_ok
  14. lb_ok = ha.Open(ls_ip, ll_port)
  15. if not lb_ok then
  16.     MessageBox('连接失败', '请检查 Redis 服务')
  17.     destroy ha
  18.     return
  19. end if
  20. ulong lul_sub
  21. lul_sub = ha.Subscribe(ls_channel, true)   // switch=true 开启订阅,返回订阅 id
  22. MessageBox('订阅已提交', 'channel=' + ls_channel + ' id=' + string(lul_sub))
  23. // 收到推送时,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():断开。

  1. // 前置:已部署 Kafka(单 broker,localhost:9092);窗口按钮 clicked 事件贴入下方代码运行。
  2. // 示例输入:Kafka broker 地址(多个用逗号分隔)
  3. string ls_brokers
  4. ls_brokers = 'localhost:9092'
  5. // 示例输入:连接超时(毫秒)
  6. long ll_timeout
  7. ll_timeout = 5000
  8. // 示例输入:目标主题
  9. string ls_topic
  10. ls_topic = 'pb_events'
  11. // 示例输入:要发送的消息
  12. string ls_msg
  13. ls_msg = 'user_login|1001'
  14. uo_kafka_producer kp
  15. kp = create uo_kafka_producer
  16. kp.ProducerConnect(ls_brokers, ll_timeout)
  17. int li_rtn
  18. li_rtn = kp.ProducerSend(ls_topic, ls_msg)   // 也可传 blob 发二进制
  19. if li_rtn = 0 then
  20.     MessageBox('发送成功', 'topic=' + ls_topic)
  21. else
  22.     MessageBox('发送失败', kp.ProducerGetError())
  23. end if
  24. kp.ProducerDisconnect()
  25. 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/GetConsumerGetError()

前置:在窗口(或 NVO)声明实例变量 iuc,并编写 iucue_message 事件;启动后持续 Poll 即可收到消息。
  1. // 1) 窗口/对象声明实例变量(实例作用域)
  2. uo_kafka_consumer iuc
  3. // 2) iuc 的 ue_message 事件(收到消息时触发):status 为状态串,data 为消息体(blob)
  4. string ls_status
  5. ls_status = status
  6. // data 是 blob,按业务解码(例如转字符串后处理)
  7. MessageBox('收到消息', ls_status)
  8. // 3) 启动消费(按钮 clicked 事件)
  9. // 示例输入:broker 地址
  10. string ls_brokers
  11. ls_brokers = 'localhost:9092'
  12. // 示例输入:连接超时(毫秒)
  13. long ll_timeout
  14. ll_timeout = 5000
  15. // 示例输入:订阅主题
  16. string ls_topic
  17. ls_topic = 'pb_events'
  18. iuc = create uo_kafka_consumer
  19. iuc.ConsumerConnect(ls_brokers, ll_timeout)
  20. int li_sub
  21. li_sub = iuc.ConsumerSubscribe(ls_topic)
  22. if li_sub = 0 then
  23.     MessageBox('订阅成功', 'topic=' + ls_topic)
  24. else
  25.     MessageBox('订阅失败', iuc.ConsumerGetError())
  26. end if
  27. // 接着在窗口 timer 事件里循环调用 iuc.ConsumerPoll(1000) 驱动消息分发
复制代码

窗口 timer 事件(先 timer(1) 让窗口每秒触发一次):
  1. // 窗口 timer 事件:周期性拉取,消息经 ue_message 回调
  2. boolean lb_more
  3. lb_more = iuc.ConsumerPoll(1000)
  4. // lb_more=true 表示还有积压消息,可继续 Poll;false 表示本次已空
复制代码

十、边界情况与常见坑


  • Redis 无密码Openauth 必须传空串 '',传 null 或省略会连不上。
  • BRPOP 是阻塞的timeout=0 会一直阻塞当前调用线程。生产环境务必设合理超时(如 5~30 秒),或把消费放到 uo_thread 后台线程,避免卡 UI。
  • reply 数组下标BRPOP 返回 [队列名, 值],取消息体用 reply_array[2](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 老系统就能以很低的成本获得"异步 + 削峰 + 解耦"的能力,而不必重写架构。
共享共进共赢
Sharing And Win-win Results
SYBASEBBS - 免责申明1、欢迎访问“SYBASEBBS.COM”,本文内容及相关资源来源于网络,版权归版权方所有!本站原创内容版权归本站所有,请勿转载!
2、本文内容仅代表作者观点,不代表本站立场,作者自负,本站资源仅供学习研究,请勿非法使用,否则后果自负!请下载后24小时内删除!
3、本文内容,包括但不限于源码、文字、图片等,仅供参考。本站不对其安全性,正确性等作出保证。但本站会尽量审核会员发表的内容。
4、如本帖侵犯到任何版权问题,请立即告知本站 ,本站将及时删除并致以最深的歉意!客服邮箱:admin@sybasebbs.com
您需要登录后才可以回帖 登录 | 站点注册

本版积分规则

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

Mail To:Admin@SybaseBbs.com

客服微信:18669893686
挂谷猜想 · 探索

QQ|Archiver|PowerBuilder(PB)BBS社区 ( 鲁ICP备2021027222号-1 )

GMT+8, 2026-9-2 14:07 , Processed in 0.025823 second(s), 9 queries , MemCached On.

Powered by Discuz! X3.5

© 2001-2026 Discuz! Team.

快速回复 返回顶部 返回列表