08-缓存与消息 common-cache-mq¶
📄 创建: ZCode AI 2026-08-14 · 修改: JimWb 2026-08-14
common-cache(12 个类)提供基于 Jedis 的 Redis 缓存基础服务和基于配置表的缓存自动加载框架。bams-mq(4 个类)提供 Kafka 消息队列封装。
Part 1: common-cache(Redis 缓存)¶
一、RedisCacheManager —— 缓存基础服务¶
cn.cnbm.bams.cache.service.RedisCacheManager(抽象类)
基于 Jedis 连接池的 Redis CRUD 服务。所有方法都是静态方法,值用 GsonUtils JSON 序列化。
连接池配置¶
| 参数 | 默认值 |
|---|---|
| maxIdle | 200 |
| maxTotal | 200 |
| minIdle | 20 |
| maxWait | 18000ms |
| testOnBorrow / testOnReturn | true |
| testWhileIdle | true |
静态 CRUD 方法¶
| 方法签名 | 功能 |
|---|---|
setObject(String key, Object value) |
设置(永久) |
setObject(String key, Object value, int expiretime) |
设置(秒级过期) |
setObject(String key, String field, Object value) |
Hash hset |
getObject(String key) |
取 Object |
<T> T getObject(String key, Class<T> clazz) |
取指定类型 |
getObject(String prefix, String key) |
按前缀+key 取 |
delKey(String key) |
删除 |
hdel(String key, String field) |
Hash 删 |
deleteByPattern(String pattern) |
按模式批量删 |
exists(String key) |
是否存在 |
expire(String key, int seconds) |
设过期 |
Hash 操作¶
| 方法 | 功能 |
|---|---|
setObjectByHSet(key, field, value) |
Hash 写入 |
getObjectByHGet(key, field) |
Hash 单字段取 |
getAllObjectByHGet(key) |
Hash 全部取 |
getObjectByHGetLikes(key, fieldPattern) |
Hash 模糊取 |
模糊查询¶
| 方法 | 功能 |
|---|---|
redisKeys(String pattern) |
按模式查 key |
redisKeysMap(String pattern) |
按模式查 key+value |
redisKeysAllMap(String pattern) |
全量 |
使用示例¶
// 基本读写
RedisCacheManager.setObject("user:1001", userObj);
RedisCacheManager.setObject("token:abc", tokenStr, 3600); // 1小时过期
User user = RedisCacheManager.getObject("user:1001", User.class);
// Hash
RedisCacheManager.setObjectByHSet("dict:status", "1", "启用");
String val = (String) RedisCacheManager.getObjectByHGet("dict:status", "1");
// 删除
RedisCacheManager.delKey("user:1001");
RedisCacheManager.deleteByPattern("session:*"); // 删所有 session: 开头的 key
二、DbService —— 配置驱动的缓存自动加载¶
cn.cnbm.bams.cache.service.DbService
基于 bd_cache 配置表,启动时自动把数据库数据加载到 Redis。支持 4 种保存模式。
bd_cache 配置表¶
| 字段 | 说明 |
|---|---|
tableName |
数据源表名 |
tableKey |
Redis key 前缀 |
keyCreate |
分组字段(决定 key 的构成) |
cacheValue |
要缓存的列 |
saveMode |
保存模式:array / object / tree / hset |
sqlWhere |
WHERE 条件 |
ext1 |
扩展(值为 "class" 时走本地方法) |
expiretime |
过期时间(秒,默认 7 天) |
isDel |
是否禁用 |
4 种保存模式¶
| 模式 | 说明 | Key 格式 | 值 |
|---|---|---|---|
array |
按字段分组 | {tableKey}:{分组值} |
List |
object |
按字段为 key | {tableKey}:{值} |
单个值 |
tree |
构建树结构 | Hash | parentCode→childs 树 |
hset |
多字段拼接 | Hash | field=多字段拼接,值 JSON |
initBaseCache 流程¶
1. 读 bd_cache 表(is_del=0)
2. 对每条配置:
├─ ext1="class" → 反射调本地方法
└─ 否则执行 select {cacheValue} from {tableName} [WHERE {sqlWhere}]
3. 按 saveMode 分组/构建
4. deleteByPattern 清旧数据
5. 写入 Redis(列名下划线转驼峰)
6. 更新 bd_cache.last_refresh_time
手动刷新¶
三、配置项¶
spring:
data:
redis:
host: 127.0.0.1
port: 6379
password:
database: 3
time-out: 2000
project:
cache:
init-cache: true # 启动时加载缓存(默认 true)
init-taos: true # 初始化 TDengine
配置类¶
| 类 | 说明 |
|---|---|
InitRedisService |
@Component @Order(1),@PostConstruct 触发缓存初始化 |
RedisCacheProperties |
Redis 连接参数 + appId/appUrl/cacheClazz |
SettingsProperties |
project.cache.init-cache / init-taos |
四、缓存 Key 前缀常量¶
cn.cnbm.bams.cache.constant.CacheProperty
| 常量 | 前缀 | 用途 |
|---|---|---|
bd_device: |
设备 | |
bd_frame_location: |
框架位置 | |
bd_point: |
测点 | |
bd_dictionary: |
字典 | |
real_point_data |
实时点位 | |
realtime_data: |
实时数据 | |
location_type: |
位置类型 | |
pub_datasource: |
数据源 |
Part 2: bams-mq(Kafka 消息队列)¶
一、KafkaProducer —— 生产者¶
cn.cnbm.bams.mq.service.KafkaProducer(@Service)
| 方法签名 | 功能 |
|---|---|
CompletableFuture<SendResult<String,String>> sendMessage(String message) |
发到默认 topic |
CompletableFuture<SendResult<String,String>> sendMessage(String topic, String message) |
发到指定 topic |
void flush() |
刷新 |
@Autowired
private KafkaProducer kafkaProducer;
// 发送
kafkaProducer.sendMessage("send-data", jsonData);
// 或用默认 topic(project.kafka.send-topic)
kafkaProducer.sendMessage(jsonData);
二、KafkaConsumer —— 消费者¶
cn.cnbm.bams.mq.service.KafkaConsumer(@Service)
必须配置
project.kafka.consumer-enabled=true才启用消费者。
特点:
- 用 topicPattern(正则匹配,默认 bams_.*)订阅,不是固定 topic
- 批量消费 List<ConsumerRecord>,手动 ack.acknowledge()
- 收到消息后调 dataProcessingService.accept(record)
三、IDataProcessingService / BaseDataProcessingService —— 消息处理¶
IDataProcessingService(接口)¶
public interface IDataProcessingService extends Runnable, Consumer<Object> {
// 继承 Consumer<Object>,实现 accept(Object) 处理消息
}
BaseDataProcessingService(抽象类,推荐继承)¶
消息处理模板基类:
@Service
public class MyMessageHandler extends BaseDataProcessingService {
@Override
protected void perHandleMessage(JSONObject dataObject) {
// 预处理
String businessCode = dataObject.getString("businessCode");
// ...
}
@Override
protected void doHandleMessage(JSONObject unit) {
// 主处理
String operateTag = unit.getString("operateTag"); // I/U/D
// ...
}
}
BaseDataProcessingHandler 流程¶
handleMessage(message)
├─ String 类型 → Gson 解析为 JSONObject
├─ List<ConsumerRecord> 类型 → 遍历,每个 value Gson 解析
├─ perHandleMessage(jsonObject) // 预处理(抽象)
├─ doHandleMessage(jsonObject) // 主处理(抽象)
└─ 捕获 BussinessException → 记录错误日志 + 异常数据
四、配置项¶
spring:
kafka:
bootstrap-servers: 192.168.1.1:9092,192.168.1.2:9092,192.168.1.3:9092
consumer:
group-id: group1
project:
kafka:
consumer-enabled: true # 必须为 true 才启用消费者
send-topic: send-data # 生产默认 topic
receive-topic-pattern: bams_.* # 消费 topic 正则