消息通道订阅(拉取模式)
消息通道订阅事件对接
1. 简介
消息通道供开发者主动拉取平台实时消息,需自行轮询消费。与平台主动推送(推送模式)消息方案互为补充。
- 消息在通道中最多保留 2 天;消息中图片的有效期为7天,需长期保存请自行转存。
msgBody结构与消息推送载荷一致,可按msgType复用解析逻辑(详见 事件消息类型定义)。
2. 开通与约束
开通: 提交技术工单或联系我们申请开通账号的消息通道订阅功能权限,约 5 分钟后生效。权限开通后,新产生的消息会进入消息通道。
| 约束 | 说明 |
|---|---|
| 消费组 | 固定 3 个:group1 / group2 / group3,数据相同、偏移量独立 |
| 并发 | 同一 appId + group 同时仅允许一个拉取或提交请求 |
| 单次条数 | 最大 500 条 |
| 部署 | 每个 appId + group 建议单实例、单线程消费 |
3. 消费流程与偏移量
获取 accessToken → 循环 pullMessages → 处理 messages
├─ autoCommit=true → pull 成功即提交本轮位点,无需 commitOffset
└─ autoCommit=false → 处理成功后调用 commitOffset
若 hasMore=true,间隔 ≥ 1 秒继续拉取
偏移量提交策略
| 模式 | 行为 | 适用场景 |
|---|---|---|
autoCommit=true | 本次 pullMessages 成功后,服务端立即提交本轮位点 | 联调、快速验证;处理失败可能丢消息 |
autoCommit=false | 位点进入 pending,业务处理成功后调用 commitOffset | 生产环境推荐 |
注意:
autoCommit=true时位点在 pull 返回后即已推进,与业务是否处理完成无关;须异步处理并接受 at-most-once 语义。autoCommit=false时,无 pending 位点调用commitOffset返回MC1003。- 建议用
msgId做幂等去重。
Java SDK: 负责轮询拉取并回调 IMessageHandler。不会调用 commitOffset HTTP 接口;autoCommit=true 时位点在 pull 请求中由服务端提交,autoCommit=false 时须在 Handler 内调用 InboxConsumer.commitOffset()(内置 token 失效重试)。SDK 内部已串行化 pull 与 commit 请求;手动模式建议在回调线程同步处理并提交。
消息消费示例Demo及SDK:inbox-consumer-sdk
4. HTTP API
4.1 pullMessages:查询消息列表,消费报警事件
请求地址
https://openapi.lechange.cn/openapi/pullMessages
传入参数说明
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
token | String | 是 | accessToken |
group | String | 是 | 消费组,可填group1 / group2 / group3 |
limit | Integer | 否 | 单次最大条数,默认 100,最大 500 |
autoCommit | Boolean | 否 | 默认 false;true 时 pull 成功后立即提交本轮位点 |
样例输入
{
"system":{
"ver":"1.0",
"appId":"lcdxxxxxxxxx",
"sign":"85384489874412a2e9aa7ac3a8b08309",
"time":1603350639,
"nonce":"350f4c7a-a352-46cf-ae03-f4f526f5d99f"
},
"id":"b8f45dbb-79ad-4a7e-8b24-7eb618c5191f",
"params":{
"token":"At_00000ad9e6e87f0142eb92e207aec46a",
"group":"group1",
"limit": 500,
"autoCommit":"false"
}
}
响应 data 主要字段:
| 字段 | 说明 |
|---|---|
messages | 消息列表,元素含 msgId、msgBody,其中msgBody值可参考消息内容示例 |
count | 本次条数 |
hasMore | 是否可能仍有更多消息 |
offsetReset | 是否因超 2 天保留期触发位点重置(部分历史消息可能丢失) |
4.2 commitOffset:提交偏移量
仅在请求pullMessages接口是传
autoCommit=false且本批消息处理成功后调用。
请求地址
https://openapi.lechange.cn/openapi/commitOffset
传入参数说明
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
token | String | 是 | accessToken |
group | String | 是 | 消费组,可填group1 / group2 / group3 |
样例输入
{
"system":{
"ver":"1.0",
"appId":"lcdxxxxxxxxx",
"sign":"85384489874412a2e9aa7ac3a8b08309",
"time":1603350639,
"nonce":"350f4c7a-a352-46cf-ae03-f4f526f5d99f"
},
"id":"b8f45dbb-79ad-4a7e-8b24-7eb618c5191f",
"params":{
"token":"At_00000ad9e6e87f0142eb92e207aec46a",
"group":"group1"
}
}
返回data字段说明
无data数据返回
样例输出
{
"id":"b8f45dbb-79ad-4a7e-8b24-7eb618c5191f",
"result":{
"code":"0",
"msg":"操作成功"
}
}
4.3 消息体
消息类型见 事件消息类型定义。 消息体(对应 pullMessages 接口返回的 msgBody 字段)见 事件消息格式定义。
5. API接口错误码
| 错误码 | 说明 | 处理建议 |
|---|---|---|
0 | 成功 | — |
OP1026 | 请求过于频繁 | 同一个appId同一个group并发调用会报OP1026,请避免多实例并发 |
MC1001 | 通道未开通 | 消息通道权限未开通,请联系商务开通并等待生效 |
MC1002 | 消费组非法 | 仅支持使用 group1/group2/group3 |
MC1003 | 无可提交的偏移量 | 先 pull 产生 pending,处理后再 commit |
网络超时等临时异常可重试;参数类错误请修正后重试。
6. 常见问题
拉取不到消息: 确认通道已开通且生效;确认该时段有新消息;hasMore=false 且 count=0 表示已追上最新;
重复消费: 手动模式未 commit 或自动模式处理过慢导致重复拉取;用 msgId 幂等。生产环境请用 autoCommit=false 并在处理成功后 commitOffset。
**offsetReset=true:** 消息超过 2 天保留期被清理,位点已重置。
7. Java SDK 集成
消息消费示例Demo及SDK:inbox-consumer-sdk
SDK(JDK 8+)封装鉴权、签名与轮询拉取,提供 commitOffset() 供手动提交使用。
7.1 获取 SDK
环境要求: JDK 8+
SDK 以 fat jar 方式提供:inbox-consumer-sdk-1.0.0.jar。
将 jar 放入业务项目的 lib/ 目录,按以下方式引入:
- 普通 Java 项目: 将
inbox-consumer-sdk-1.0.0.jar加入 classpath 即可。 - Maven 项目: 在
pom.xml中以systemscope 引用本地 jar:
<properties>
<sdk.version>1.0.0</sdk.version>
</properties>
<dependency>
<groupId>com.imou.open.local</groupId>
<artifactId>inbox-consumer-sdk-bundle</artifactId>
<version>${sdk.version}</version>
<scope>system</scope>
<systemPath>${project.basedir}/lib/inbox-consumer-sdk-${sdk.version}.jar</systemPath>
</dependency>
7.2 配置方式
在代码中声明参数,通过 InboxConsumer.builder() 传入(与 inbox-consumer-demo 的 Main.java 一致):
| 配置项 | 必填 | SDK 默认 | 说明 |
|---|---|---|---|
appId / appSecret | 是 | — | 开放平台应用凭证 |
gatewayUrl | 否 | https://openapi.lechange.cn/openapi | OpenAPI 根路径 |
group | 否 | group1 | 消费组:group1 / group2 / group3 |
limit | 否 | 100 | 单次拉取条数(1~500) |
autoCommit | 否 | true | 联调可保持默认(pull 时服务端提交位点);生产请设为 false 并手动 commit |
intervalMs | 否 | 1000 | 拉取间隔(毫秒,≥1000) |
7.3 接入示例
IMessageHandler.onMessages 入参为 InboxPullResult,可读取单次 pull 的元数据:
| 字段 | 方法 | 说明 |
|---|---|---|
messages | getMessages() | 本次拉取的消息列表 |
count | getCount() | 本次拉取条数 |
hasMore | isHasMore() | 是否可能仍有更多消息 |
offsetReset | isOffsetReset() | 是否因超过 2 天保留期触发位点重置 |
main 中声明配置变量,再传给 builder() 与 Handler(生产环境将 autoCommit 设为 false):
import com.imou.open.inbox.sdk.*;
import java.io.IOException;
public class MyApp {
public static void main(String[] args) throws Exception {
String appId = "你的开发者appId";
String appSecret = "你的开发者appSecret";
String gatewayUrl = "https://openapi.lechange.cn/openapi";
String group = "group1";
int limit = 100;
boolean autoCommit = false;
long intervalMs = 1000L;
MessageHandler handler = new MessageHandler(autoCommit);
InboxConsumer consumer = InboxConsumer.builder()
.appId(appId)
.appSecret(appSecret)
.gatewayUrl(gatewayUrl)
.group(group)
.limit(limit)
.autoCommit(autoCommit)
.intervalMs(intervalMs)
.messageHandler(handler)
.build();
handler.bind(consumer);
consumer.start();
Runtime.getRuntime().addShutdownHook(new Thread(consumer::stop));
Thread.currentThread().join();
}
static class MessageHandler implements IMessageHandler {
private final boolean autoCommit;
private InboxConsumer consumer;
MessageHandler(boolean autoCommit) {
this.autoCommit = autoCommit;
}
void bind(InboxConsumer consumer) {
this.consumer = consumer;
}
@Override
public void onMessages(InboxPullResult pullResult) {
if (pullResult.isOffsetReset()) {
// 位点已重置,超过保留期的历史消息可能无法消费
}
if (autoCommit) {
// 异步处理,避免阻塞拉取线程(示例略)
return;
}
for (InboxMessage msg : pullResult.getMessages()) {
// 按 msg.getMsgType() 处理 msg.getMsgBody()
}
// hasMore / count 可按需用于监控或背压
boolean hasMore = pullResult.isHasMore();
int count = pullResult.getCount();
try {
consumer.commitOffset();
} catch (IOException | InterruptedException e) {
throw new IllegalStateException("commitOffset failed", e);
}
}
}
}
**autoCommit 与处理方式:**
autoCommit | 处理方式 | 偏移量 |
|---|---|---|
false | 在 onMessages 回调线程同步处理并 commitOffset() | 业务调用 commitOffset() |
true | 必须异步,避免阻塞拉取线程 | pull 成功时服务端已提交,勿再 commit |
处理失败时不要调用 commitOffset(),否则可能丢失未处理完的消息。
7.4 Demo 快速验证
提供 inbox-consumer-demo 示例工程(已包含 lib/inbox-consumer-sdk-1.0.0.jar)。
- 配置: 编辑
inbox-consumer-demo/src/main/java/com/imou/open/inbox/demo/Main.java,修改main方法开头的变量:
String appId = "你的开发者appId";
String appSecret = "你的开发者appSecret";
String gatewayUrl = "https://openapi.lechange.cn/openapi";
String group = "group1";
int limit = 100;
boolean autoCommit = false; // 联调可改为 true
long intervalMs = 1000L;
- 运行:
- IntelliJ: 用 IDEA 打开
inbox-consumer-demo目录,执行 Maven Reload Project 后运行Main。 - 命令行:
java -jar inbox-consumer-demo-1.0.0.jar
Demo 中 autoCommit=true 时异步处理消息;autoCommit=false 时同步处理并在成功后调用 consumer.commitOffset()。