Siddhi对接Google Pub/Sub如何仅在事件执行成功时确认消息
Siddhi 对接 Google Pub/Sub 实现处理成功才确认消息的可行方案
核心是替换默认的接收即自动确认机制,将消息确认动作和Siddhi流处理的最终结果绑定,具体操作步骤如下:
- 调整输入源配置,关闭自动确认
在@source配置块中显式添加auto.ack="false"参数,直接关闭插件收到消息就立刻ack的默认逻辑。此时未被显式确认的消息会在配置的ack超时时间到期后,由Pub/Sub服务端自动重新投递给订阅者,从根源上避免处理失败但消息被误确认的问题。
基础配置示例:@source(type = 'googlepubsub', project.id = 'gcp-project-xxx', topic.id = 'event-topic', subscription.id = 'event-sub', auto.ack = "false", @map(type = 'json')) define stream EventInput (eventId string, payload string); - 在处理流末端绑定手动确认逻辑
关闭自动ack后,插件会在每条流入的事件中隐式注入_pubsub_ack_id标识字段,你只需要把ack触发逻辑放在所有业务处理流程的最末端即可:- 若使用2.x及以上版本的扩展,直接在所有处理逻辑(数据校验、转换、下游sink写入等)全部执行成功的节点,调用内置函数
pubsub:ack(_pubsub_ack_id)完成确认。只要流程中途抛出异常、执行中断,没有走到该函数调用步骤,对应消息就不会被确认,后续会自动重投。 - 若处理流程包含窗口聚合、异步输出等算子,必须等所有算子返回成功结果后再触发ack,禁止在流程中途提前确认。
- 若使用2.x及以上版本的扩展,直接在所有处理逻辑(数据校验、转换、下游sink写入等)全部执行成功的节点,调用内置函数
- 配套容错配置
- 给Pub/Sub订阅设置合理的ack超时时间(
ackDeadlineSeconds),数值要略大于单条消息的最大处理耗时,避免消息还在处理中就因超时被重复投递。 - 给订阅配置死信主题,设置最大重投次数(通常设5-10次即可),多次处理失败的异常消息会自动转入死信主题,不会阻塞正常消费流程。
- 给Pub/Sub订阅设置合理的ack超时时间(
注意:1.0.x及更早版本的siddhi-io-googlepubsub扩展硬编码了接收即ack的逻辑,没有暴露手动确认的配置入口,如果当前使用的是老版本,需要先将扩展升级到1.1.0以上的稳定版本再做上述配置。
内容的提问来源于stack exchange,提问作者Sebastian Wagner
相关产品推荐
相关产品推荐

