暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

分布式实战:缓存数据生产服务之订阅通知

原创 Oracle 2023-06-12
176


本文首发于Ressmix个人站点:https://www.tpvlog.com

我们前面章节已经搭建完了Zookeeper和Kafka集群,本章我们可以基于Kafka,完成缓存数据生产服务的剩余功能,我将带领大家实现缓存数据生产服务对Kafka的商品信息Topic的订阅和消费,获取商品基本信息和店铺信息的变更通知,然后查询商品基本信息和店铺信息,最后更新本地缓存和Redis分布式缓存。

本章的应用代码存放在Gitee:5.订阅Kafka(https://gitee.com/ressmix/epay/tree/master/5.%E8%AE%A2%E9%98%85Kafka)的epay-cache
项目下。

一、Kafka消费者

缓存数据生产服务将作为消费者从Kafka的topic中消费数据(通知),所以我们需要在缓存数据生产服务中引入kafka客户端:

1.1 Kafka配置

首先,我们需要在应用的maven依赖中引入kafka依赖:

1<dependency>
2    <groupId>org.springframework.kafka</groupId>
3    <artifactId>spring-kafka</artifactId>
4    <version>2.5.1.RELEASE</version>
5</dependency>

因为Spring为Kafka提供了很好的支持,并提供了与原生Kafka Java客户端一起使用的抽象层,所以我们这里直接引用spring-kafka就可以了。

关于更多Spring集成Kafka的资料,可以参考Spring官方文档:https://spring.io/projects/spring-kafka/。

接着,我们需要在application.properties
配置上加上Kafka的配置项:

 1# 整体配置
2spring.kafka.bootstrap-servers=192.168.0.107:9092,192.168.0.109:9092,192.168.0.110:9092
3spring.kafka.listener.missing-topics-fatal=false
4# 消费者配置
5spring.kafka.consumer.enable-auto-commit=false
6spring.kafka.consumer.auto-commit-interval=1000ms
7spring.kafka.consumer.auto-offset-reset=earliest
8spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
9spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
10# 生产者配置
11spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
12spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer

上述的配置都是针对Kafka的全局配置,都是一些Kafka常用的参数,读者可以参考Spring官方文档详细了解。

最后,再增加一个Kafka配置类:

 1@Configuration
2@EnableKafka
3@EnableConfigurationProperties(KafkaProperties.class)
4public class KafkaConfig {
5
6    @Autowired
7    private KafkaProperties kafkaProperties;
8
9    @Bean
10    KafkaTemplate<String, Object> createKafkaTemplate() {
11        Map<String, Object> props = new HashMap<String, Object>();
12        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
13        props.put(ProducerConfig.RETRIES_CONFIG, 3);
14        props.put(ProducerConfig.LINGER_MS_CONFIG, 1);
15        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, kafkaProperties.getProducer().getKeySerializer());
16        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, kafkaProperties.getProducer().getValueSerializer());
17
18        KafkaTemplate<String, Object> template = new KafkaTemplate<String, Object>(new DefaultKafkaProducerFactory(props));
19        return template;
20    }
21}

1.2 订阅通知

接着,我们需要创建通知监听类,用于监听Kafka中的商品信息和店铺信息的变动通知,我这里用了两个Topic:“product-topic”和"shop-topic",分别对应商品基本信息服务和店铺服务:

 1@Component
2public class CacheIntegrationConsumer {
3
4    @Autowired
5    private CacheService cacheService;
6
7    private static final Logger LOG = LoggerFactory.getLogger(CacheIntegrationConsumer.class);
8
9    @KafkaListener(topics = {"product-topic"}, groupId = "product-group")
10    public void consumeProduct(ConsumerRecord<Integer, String> record) {
11        LOG.info("接收到商品服务通知: {}", record);
12        String data = record.value();
13
14        // 提取出商品id
15        Long productId = Long.valueOf(data);
16
17        // 调用商品基本信息服务的接口,获取最新数据
18        // 生产环境一般RPC调用,这里直接注释模拟
19        String productInfoJSON = "{\"id\": 1,\"productId\": 1000, \"name\": \"iphone7手机\", \"price\": 5599, \"pictureList\":\"a.jpg,b.jpg\", \"specification\": \"iphone7的规格\", \"service\": \"iphone7的售后服务\", \"color\": \"红色,白色,黑色\", \"size\": \"5.5\", \"shopId\": 15, \"modifiedTime\": \"2020-01-01 12:00:00\"}";
20        ProductInfo productInfo = JSONObject.parseObject(productInfoJSON, ProductInfo.class);
21
22        // TODO:重建本地缓存和Redis缓存
23        cacheService.rebulidCache(productInfo);
24    }
25
26    @KafkaListener(topics = {"shop-topic"}, groupId = "shop-group")
27    public void consumeShop(ConsumerRecord<Integer, String> record) {
28        LOG.info("接收到店铺服务通知: {}", record);
29        JSONObject jsonObject = JSONObject.parseObject(record.value());
30
31        // 提取出店铺id
32        Long shopId = jsonObject.getLong("shopId");
33
34        // 调用店铺服务的接口
35        // 生产环境一般RPC调用,这里直接注释模拟
36        String shopInfoJSON = "{\"id\": 1,\"shopId\": 15, \"name\": \"小王的手机店\", \"level\": 5, \"goodCommentRate\":0.99, \"modifiedTime\": \"2020-01-01 13:00:00\"}";
37        ShopInfo shopInfo = JSONObject.parseObject(shopInfoJSON, ShopInfo.class);
38
39        // TODO:重建本地缓存和Redis缓存
40        cacheService.rebulidCache(shopInfo);
41    }
42}

上述为了简便起见,我不再去通过RPC调用商品基本信息服务和店铺服务了,只要了解这种架构和设计思路即可,因为真正业务需求千变万化。

另外,最后更新Redis缓存和本地JVM缓存的地方,调用了cacheService.rebulidCache()
方法,内部其实就是缓存的重建,我会在下一章节详细讲解。

二、测试

最后,我们启动应用,然后到服务器ressmix-dsf03上创建一个生产者,往product-topic
里面扔几条消息,模拟通知:

1[root@ressmix-dsf03 config]# /usr/local/kafka_2.12-2.5.0/bin/kafka-console-producer.sh --bootstrap-server 192.168.0.107:9092,192.168.0.109:9092,192.168.0.110:9092 --topic product-topic 
2>1002012300102
3>

可以看到,缓存数据生产服务成功订阅到了通知,并最终更新了JVM和Redis缓存:

12020-05-28 23:15:43 [com.tpvlog.epay.cache.kafka.CacheIntegrationConsumer:22] INFO  - 接收到商品服务通知: ConsumerRecord(topic = product-topic, partition = 0, leaderEpoch = 0, offset = 3, CreateTime = 1590678940652, serialized key size = -1, serialized value size = 13, headers = RecordHeaders(headers = [], isReadOnly = false), key = null, value = 1002012300102)
22020-05-28 23:15:47 [com.tpvlog.epay.cache.kafka.CacheIntegrationConsumer:34] INFO  - 本地缓存的商品信息:com.tpvlog.epay.cache.entity.ProductInfo@322d8b01[color=<null>,id=1,name=iphone7手机,pictureList=<null>,price=5599,service=<null>,shopId=<null>,size=<null>,specification=<null>]
32020-05-28 23:15:47 [com.tpvlog.epay.cache.kafka.CacheIntegrationConsumer:38] INFO  - Reids缓存的商品信息:com.tpvlog.epay.cache.entity.ProductInfo@675263dd[color=<null>,id=1,name=iphone7手机,pictureList=<null>,price=5599,service=<null>,shopId=<null>,size=<null>,specification=<null>]

复制

三、总结

本章,我在epay-cache
这个缓存数据生产服务项目下编写了Kafka的客户端使用代码,缓存数据生产服务将作为消费者客户端监听topic的消息,然后调用源服务的数据接口,最后更新Redis缓存和本地JVM缓存。

事实上,在高并发场景下,我们针对缓存的重建还需要做一些优化,否则可能会出现缓存风暴和脏写等问题。下一章,我将会详细讲解缓存重建的方案。

「喜欢这篇文章,您的关注和赞赏是给作者最好的鼓励」
关注作者
【版权声明】本文为墨天轮用户原创内容,转载时必须标注文章的来源(墨天轮),文章链接,文章作者等基本信息,否则作者和墨天轮有权追究责任。如果您发现墨天轮中有涉嫌抄袭或者侵权的内容,欢迎发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论