修改Spring Kafka监听器适配文件切换目标staging表的方法
复用Kafka代码库,将数据导入我方Staging表而非对方表
我想复用其他团队的Kafka应用/代码库,消费相同的Kafka数据,但要把数据加载到我方新的staging表,而不是对方的目标表。对方代码在“Messages”文件夹下有多个Kafka监听器适配.java文件,每个对应消费不同类型的数据,里面包含大量导入语句、公共类和执行数据插入的try块等。下面是需要修改的核心内容:
一、调整Kafka监听器的GroupID与ID
首先得改@KafkaListener注解里的groupId和id参数,同时在配置文件里新增对应配置项。
原注解代码:
@KafkaListener(topics = "#{new java.util.HashMap(${dataflow.consumer.props.assets.crypto}).keySet()}", groupId = "${dataflow.consumer.assets.crypto.group}",containerFactory = "cryptoListenerContainerFactory", id = "${dataflow.consumer.assets.crypto.group}")
修改步骤:
- 把注解里
groupId和id对应的配置占位符换成我方的配置键,比如${dataflow.consumer.assets.crypto.my-team-group} - 在配置文件(比如application.properties)里加我方的GroupID配置:
# 基础Kafka消费者配置 spring.kafka.consumer.group-id=我方新GroupID spring.kafka.listener.id=我方新ListenerID # 对应注解里的自定义占位符配置 dataflow.consumer.assets.crypto.my-team-group=我方新GroupID
二、替换数据插入的Repository关联
数据最终是通过Repository的storecryptoCollection方法插入表的,这里得确保Repository指向我方的staging表实现:
原插入代码块:
try { List<FundEntity> cryptoFilteredList = cryptoConsumerRecordList.stream() .filter(cryptoConsumerRecord -> recordIsNotFilterable(cryptoConsumerRecord, FundEntity.class)) .map(ConsumerRecord::value) .collect(Collectors.toList()); if (!cryptoFilteredList.isEmpty()) { this.Repository.storecryptoCollection(cryptoFilteredList, topic, firstOffset, partition, cryptoGroupId); kafkaDetailsSummary.setStatus(KafkaMessageSummary.KafkaStatus.PROCESSED); log.info(LogHelper.generateCryptoReconMessage(cryptoFilteredList), keyValue(KAFKA_DETAILS, kafkaDetailsSummary)); }
修改要点:
- 检查
this.Repository是不是对方团队的实例,如果是,得换成我方对应staging表的Repository实现类:- 如果是依赖注入(比如@Autowired),修改注入的Repository类型,或者在配置里指定我方Repository的Bean
- 查看
storecryptoCollection方法的实现:- 如果方法里硬编码了目标表名,直接改成我方staging表的名称
- 如果是通过配置文件指定表名,在配置里新增我方表名的配置项,再修改Repository里读取表名的占位符
三、额外需要检查的地方
- 实体类映射:比如
FundEntity如果用JPA/Hibernate的@Table注解指定了表名,得把name属性改成我方staging表的名字 - 权限验证:确保我方的Kafka消费者GroupID有权限消费对应Topic,同时我方数据库账号有权限写入目标staging表
内容的提问来源于stack exchange,提问作者ncisco3200
相关产品推荐
相关产品推荐

