Spring Kafka存数据到Cassandra如何获取DB响应判断存储是否成功
Cassandra写入结果判断+逻辑调整方案
原理说明
Spring Data Cassandra 提供的 Repository 接口的 save() 方法有明确的成功/失败标识:
- 操作成功:方法返回已持久化的User实体对象,无异常抛出
- 操作失败:直接抛出
DataAccessException系列异常,覆盖连接故障、数据校验不通过、写入一致性不满足、主键冲突等所有入库失败场景
你只需要调整执行顺序,把入库逻辑放到邮件发送前,通过异常捕获即可控制邮件触发逻辑。
调整后完整代码
package com.example.demo.consumer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.dao.DataAccessException; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; import com.example.demo.User; import com.example.demo.UserRepository; @Service public class MessageConsumer { @Autowired UserRepository userrepo; @Autowired EmailSenderService senderservice; @KafkaListener(topics="k2-topic",groupId="mygroup2") public void consumeFromTopic(String message) { System.out.println("consumed message "+message); String vin = message.substring(0, 17); String verified = message.substring(17, 18); int sp = Integer.parseInt(message.substring(18,21)); String alert = message.substring(21,22); char ch = alert.charAt(0); String timeStamp = message.substring(22); String bd = "Hai There !\n Your CAR VIN NO IS : "+vin+"\nYou're crossing your SPEED LIMIT : "+sp+"\nPlease Drive slowly\nSafe driving saves life."; User obj = new User(vin,verified,sp,alert,timeStamp); try { // 先执行入库操作 userrepo.save(obj); // 入库成功,再判断是否需要发送告警邮件 if(ch == 'y') { senderservice.sendEmail(bd); } } catch (DataAccessException e) { // 入库失败,打印日志,不触发邮件 System.err.println("用户数据入库Cassandra失败,vin:"+ vin + ",错误信息:" + e.getMessage()); // 此处可根据业务需要添加重试、死信队列转发等逻辑 } } }
可选增强方案
如果需要更严谨的一致性校验,可调整配置:
- 在Cassandra配置文件中调高写入一致性等级,比如设置为
QUORUM,保证多数节点写入成功才判定操作成功 - 若业务要求强校验,可在save执行成功后新增一次主键查询,确认数据确实存在(会额外增加读请求开销,非必要不建议使用)
内容的提问来源于stack exchange,提问作者Niranjan Vasadi
相关产品推荐
相关产品推荐

