马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。
您需要 登录 才可以下载或查看,没有账号?站点注册
×
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[2] 才是真正的消息体。
前置:同上,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) // 取出数组:[1]=队列名, [2]=消息体
- string ls_val
- reply_array[2].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[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 老系统就能以很低的成本获得"异步 + 削峰 + 解耦"的能力,而不必重写架构。 |