-
摘要政企、能源、金融、军工等信创行业的数据中心与企业机房,普遍面临资产管理硬件适配难、系统不兼容、数据不合规等痛点。传统进口RFID设备、商用通用采集硬件,无法适配鲲鹏ARM架构、openEuler国产操作系统,且在算力机房强电磁、高密度金属设备场景下识别稳定性差,难以满足等保2.0与信创验收标准。本文以国产化MC-RFID硬件体系为核心,结合华为云IoT、GaussDB、边缘计算能力,从硬件选型、边缘适配、云端对接、数据合规、项目落地全维度,讲解信创场景专属RFID资产管理解决方案,彻底解决传统资管账实不符、盘点低效、国产化适配缺失、审计溯源难等问题。本文全套国产化RFID硬件资管落地方案由首码信息深耕信创物联网领域多年打磨优化,已在多个省级能源、政务算力中心落地验收,是信创行业标准化、可复用的资产数字化方案。一、信创场景RFID资产管理的核心硬件痛点区别于普通商用机房,信创涉密机房、国产化算力中心对底层采集硬件的兼容性、稳定性、安全性、国产化属性有着严苛要求,传统RFID资管方案的短板被无限放大,核心问题均集中在硬件层面:1. 生态适配断层:主流进口超高频RFID硬件、商用读写终端均基于x86架构开发,无法兼容华为鲲鹏ARM服务器与openEuler国产系统,软硬件适配成本极高,甚至出现完全无法接入的情况,不符合信创全栈替代要求。2. 复杂场景识别失效:GPU算力机柜、工业设备集群金属屏蔽严重,机房强电磁干扰密集,普通RFID硬件信号衰减、串读漏读问题频发,资产识别准确率不足80%,无法支撑精细化资管。3. 离线数据断层:涉密信创机房多采用物理隔离、断网运行模式,通用RFID硬件无本地缓存能力,断网后无法采集存储资产数据,联网后数据缺失,台账长期账实不符。4. 合规能力缺失:商用硬件无国密加密传输机制,采集日志无法实现不可篡改留存,无法满足等保2.0审计溯源、信创国产化验收的硬性标准。5. 资管维度单一:传统RFID硬件仅能实现资产盘点识别,无法联动机柜U位状态、设备功耗、环境温湿度数据,无法实现算力机房一体化运维管控。二、信创级全栈国产化RFID硬件体系选型针对信创场景多重适配难题,首码信息自研适配华为云鲲鹏生态的国产化MC-RFID硬件矩阵,摒弃传统超高频技术架构,采用低频磁耦合传感技术,从硬件底层实现国产化替代、抗干扰升级与合规能力补齐,全套硬件无境外技术依赖,完美适配openEuler系统与国产数据库生态。1. 国产MC-RFID磁耦合传感硬件(核心硬件)专为信创算力机房、工业涉密场景研发,替代传统进口RFID设备,依托磁耦合感应传输原理,彻底规避金属屏蔽、电磁干扰问题,支持±0.5U机柜高精度定位,可精准识别每一台设备的U位在位状态、移位变动,适配高密度GPU算力集群长期稳定运行。2. 国产化抗金属加密RFID标签全元器件国产自研,支持国密SM4加密存储资产信息,耐高温、抗老化、防脱落,适配服务器、交换机、工业设备等全品类金属固定资产。每枚标签绑定唯一EPC编码,实现一物一码、全生命周期溯源,适配信创资产合规管控要求。3. openEuler适配型边缘采集网关基于ARM鲲鹏架构深度适配,完美兼容openEuler全系国产操作系统,支持RFID硬件数据本地缓存、断网续传、脏数据预处理。可离线留存7天以上资产采集数据,机房恢复网络后自动同步至华为云,彻底解决隔离机房数据断层问题。4. 国产工业级手持读写终端适配信创运维规范,支持离线批量盘点、加密数据采集、资产快速绑定与解绑,每秒可批量识别数百枚标签,大幅提升大型机房、园区资产盘点效率,适配涉密场景无外网作业需求。三、基于华为云的国产化RFID资管整体架构依托华为云IoT、GaussDB、边缘计算、日志审计服务原生能力,搭配首码信息国产化RFID硬件集群,搭建「硬件感知-边缘处理-云端管控-业务应用」四层全栈国产化架构,全程适配信创与等保合规要求。1. 感知层:国产化MC-RFID标签、U位传感模块、手持终端、固定式读写器完成全域固定资产、机柜U位状态、设备运行数据实时采集;2.边缘层:openEuler边缘网关完成数据过滤、去重、加密、本地缓存,规避云端带宽压力,保障离线场景正常作业;3. 云端层:华为云IoT实现硬件设备统一接入、状态监控、远程运维;GaussDB承载资产主数据、台账、工单数据;LTS日志服务留存全链路操作日志,实现不可篡改审计溯源;4. 应用层:搭建国产化资产管理平台,实现资产自动盘点、U位资源可视化、移位异常告警、全生命周期追溯、合规报表自动生成。四、openEuler边缘硬件部署与云端对接实操基于首码信息落地适配经验,针对openEuler国产系统优化轻量化部署方案,适配涉密隔离机房,无需外网即可完成硬件调试与数据采集,联网后自动同步华为云IoT平台,部署简单、兼容性强。 # openEuler 系统国产化RFID硬件采集服务部署脚本 # 安装国产适配依赖组件 yum install mqtt-client sqlite -y # 启动RFID硬件数据采集与本地缓存服务 nohup python3 domestic_rfid_collect.py & # 联网自动同步硬件采集数据至华为云IoT python3 cloud_sync_service.py 边缘网关可自主完成硬件设备在线监测、无效数据过滤、资产状态校验,有效避免高密度场景下的串读、误读问题,将资产识别准确率稳定维持在99.9%,同时降低云端算力消耗。五、核心落地能力与行业价值依托华为云国产化生态与首码信息RFID硬件硬实力,方案解决信创资管行业核心痛点,实现全方位数字化、合规化升级:1. 全栈信创适配:从RFID采集硬件、openEuler边缘系统到华为云GaussDB数据库,全程国产化闭环,无境外技术依赖,顺利通过信创适配验收;2. 复杂场景精准采集:磁耦合RFID技术彻底解决算力机房金属屏蔽、电磁干扰难题,适配GPU高密度机柜、工业车间、涉密机房等极端场景;3. 合规审计全覆盖:硬件数据传输国密SM4加密,全流程操作日志不可篡改留存,满足等保2.0、金融、能源行业审计溯源标准;4. 运维效率大幅升级:替代人工纸质盘点,数万级资产盘点从数天缩短至数小时,人力成本降低90%以上,U位资源利用率提升25%以上;5. 离线在线双适配:支持物理隔离机房离线作业、联网自动同步,适配各类涉密、隔离信创场景管控需求。六、落地案例总结某省级能源大数据中心国产化GPU算力集群,全面落地首码信息国产化RFID硬件资管方案,结合华为云IoT生态完成全栈数字化改造。改造后,机房资产账实不符问题彻底清零,资产异常移位实时告警,全年无资产流失、错放问题;机房盘点效率提升90%,机柜U位资源利用率提升27%,有效节约机房扩容成本;全栈国产化架构顺利通过信创验收与等保三级测评,成为能源行业信创资产管理标杆案例。文末总结信创行业资产管理的数字化转型,核心在于底层硬件的国产化适配与场景化落地,只有兼容国产生态、抗干扰、高合规的RFID硬件体系,才能真正实现机房资产精细化、智能化、合规化管控。首码信息深耕国产化RFID硬件研发与信创资管落地,依托华为云鲲鹏、openEuler全栈国产生态,为政务、能源、金融、国企等行业提供从硬件部署、数据对接、系统搭建到合规验收的一体化资产管理解决方案,助力信创产业数字化高质量升级。标签:#国产化RFID硬件 #信创资产管理 #openEuler #华为云IoT #GaussDB #算力机房资管 #工业级RFID #国产化运维
-
kafka需要删除所有broker步骤是怎么样?是这样?然后删除主机吗?是否需要在zookeeper侧清理?服务界面的kafka是否还会存在?1.登录FusionInsight Manager。2.单击“集群 > 待操作集群的名称 > 服务 > Kafka”,进入Kafka服务状态页面。3.单击“实例”页签,进入Kafka实例页面。4.勾选待删除的实例复选框。5.选择“更多 > 删除”,执行相应的操作。6.选择“主机” > "更多 > 删除"
-
在创建布防告警任务,获取kafka信息的时候
-
ava集成kafka要在 Java 项目中集成 Apache Kafka 以实现消息的生产和消费,步骤如下:1. 引入 Maven 依赖在您的 pom.xml 文件中添加以下依赖,以包含 Kafka 客户端库:<dependencies> <!-- Kafka Clients --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>2.8.0</version> </dependency> <!-- 如果使用 Spring Boot,可添加以下依赖 --> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>2.7.0</version> </dependency> </dependencies>2. 配置 Kafka 生产者首先,设置生产者的配置属性: import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; public class KafkaProducerExample { public static void main(String[] args) { // 配置属性 Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 创建生产者 Producer<String, String> producer = new KafkaProducer<>(props); // 发送消息 for (int i = 0; i < 10; i++) { ProducerRecord<String, String> record = new ProducerRecord<>("your_topic", "key" + i, "value" + i); producer.send(record); } // 关闭生产者 producer.close(); } } 3. 配置 Kafka 消费者接下来,设置消费者的配置属性,并订阅主题以消费消息: import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class KafkaConsumerExample { public static void main(String[] args) { // 配置属性 Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "your_group_id"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 创建消费者 Consumer<String, String> consumer = new KafkaConsumer<>(props); // 订阅主题 consumer.subscribe(Collections.singletonList("your_topic")); // 持续消费消息 try { while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); records.forEach(record -> { System.out.printf("Consumed message: key = %s, value = %s, offset = %d%n", record.key(), record.value(), record.offset()); }); } } finally { // 关闭消费者 consumer.close(); } } } 4. 使用 Spring Boot 集成 Kafka如果您使用 Spring Boot,可以通过配置 KafkaTemplate(用于生产消息)和使用 @KafkaListener 注解(用于消费消息)来简化 Kafka 的集成。生产者配置: import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import java.util.HashMap; import java.util.Map; @Configuration public class KafkaProducerConfig { @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } } 使用 KafkaTemplate 发送消息: import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; @Service public class KafkaProducerService { @Autowired private KafkaTemplate<String, String> kafkaTemplate; public void sendMessage(String topic, String key, String value) { kafkaTemplate.send(topic, key, value); } } 消费者配置: import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import java.util.HashMap; import java.util.Map; @EnableKafka @Configuration public class KafkaConsumerConfig { @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "your_group_id"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } } 使用 @KafkaListener 消费消息:在 Spring Boot 中,@KafkaListener 注解用于监听指定的 Kafka 主题,并在收到消息时触发相应的方法。以下是一个基本示例: import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; @Service public class KafkaConsumerService { @KafkaListener(topics = "your_topic", groupId = "your_group_id") public void listen(String message) { System.out.println("Received message: " + message); // 在此处添加处理逻辑 } } 在上述代码中:topics:指定要监听的 Kafka 主题。groupId:指定消费者组 ID。listen 方法:当有新消息发布到指定主题时,该方法会被调用,message 参数包含消息的内容。批量消费消息如果希望一次处理多条消息,可以启用批量监听。首先,需要配置一个支持批量消费的 KafkaListenerContainerFactory: import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; @EnableKafka @Configuration public class KafkaConsumerConfig { @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory( ConsumerFactory<String, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setBatchListener(true); // 启用批量监听 return factory; } } 然后,在消费者服务中使用 @KafkaListener 注解,并指定使用上述配置的工厂: import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; import java.util.List; @Service public class KafkaBatchConsumerService { @KafkaListener( topics = "your_topic", groupId = "your_group_id", containerFactory = "kafkaListenerContainerFactory" ) public void listen(List<String> messages) { System.out.println("Received batch messages: " + messages); // 在此处添加批量处理逻辑 } } 在上述代码中:containerFactory:指定使用支持批量消费的工厂。listen 方法的参数类型为 List<String>,用于接收一批消息。控制消费者的启动和停止在某些情况下,可能需要在运行时控制 Kafka 消费者的启动和停止。可以通过 KafkaListenerEndpointRegistry 来实现: import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.listener.KafkaListenerEndpointRegistry; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; @Service public class KafkaListenerManager { @Autowired private KafkaListenerEndpointRegistry registry; // 启动监听器 public void startListener(String listenerId) { MessageListenerContainer listenerContainer = registry.getListenerContainer(listenerId); if (listenerContainer != null && !listenerContainer.isRunning()) { listenerContainer.start(); } } // 停止监听器 public void stopListener(String listenerId) { MessageListenerContainer listenerContainer = registry.getListenerContainer(listenerId); if (listenerContainer != null && listenerContainer.isRunning()) { listenerContainer.stop(); } } } 在上述代码中:startListener 方法用于启动指定的监听器。stopListener 方法用于停止指定的监听器。listenerId 对应于 @KafkaListener 注解中的 id 属性。通过这种方式,可以在应用运行时根据需要动态地控制 Kafka 消费者的行为。通过上述配置和代码示例,可以在 Spring Boot 项目中有效地集成 Kafka,实现消息的生产和消费功能。
-
问题:FI Kafka主题监控指标数据存储在哪里?需求:需要获取数据做外部系统主题流量监控
-
Kafka生产者流程是一个复杂但高效的过程,它涉及到多个步骤和组件的协同工作。以下是对Kafka生产者流程的详细分析:一、创建生产者实例配置生产者属性:生产者需要配置一系列参数以连接到Kafka集群,如bootstrap.servers(Kafka集群的地址列表)。配置序列化器(Serializer),用于将消息键(Key)和消息值(Value)序列化为字节数组,以便在网络上传输和存储。常见的序列化器有StringSerializer、ByteArraySerializer等。其他重要配置项包括acks(指定消息需要多少个副本成功写入后才认为消息发送成功)、retries(消息发送失败后的重试次数)、batch.size(控制生产者批量发送消息的大小)、linger.ms(控制生产者在发送消息之前等待更多消息加入批次的时间)等。创建KafkaProducer实例:使用配置好的属性创建一个KafkaProducer实例。这个实例将负责后续的消息发送操作。二、构建消息创建ProducerRecord:生产者通过创建ProducerRecord对象来构建要发送的消息。ProducerRecord包含了目标主题(Topic)、分区(Partition,可选)、消息键(Key,可选)和消息值(Value)。三、发送消息异步发送或同步发送:生产者可以选择异步发送或同步发送消息。异步发送可以提高吞吐量,但可能无法立即获得发送结果;同步发送则可以在发送后立即获得发送结果。在异步发送中,生产者将消息添加到缓冲区中,并异步地将缓冲区中的消息批量发送到Kafka集群。序列化消息:在发送之前,生产者会使用配置的序列化器将消息键和消息值序列化为字节数组。选择分区:Kafka根据消息键和分区策略(如轮询或哈希)选择目标分区。如果消息键为空,则使用轮询策略将消息均匀分配到各个分区。发送至Leader Broker:生产者将序列化后的消息发送到目标分区的Leader Broker。Leader Broker是分区中负责处理读写请求的Broker。四、消息确认写入本地日志文件:Leader Broker接收到消息后,将其写入本地日志文件。这是Kafka实现消息持久化的关键步骤。副本同步:Leader Broker将消息同步到从副本(Follower)Broker。从副本将消息写入本地日志文件,并向Leader发送确认。消息提交:当Leader Broker收到足够数量的从副本确认后,将消息标记为已提交(Committed)。已提交的消息对消费者可见。发送ACK给生产者:根据acks参数的设置,Leader Broker向生产者发送确认(ACK)。如果acks=all,则等待所有ISR副本(In-Sync Replicas)确认后才发送ACK。五、生产者行为调整缓冲区管理:生产者有一个内部缓冲区用于存储待发送的消息。当缓冲区满或达到发送条件时(如batch.size达到或linger.ms超时),生产者将缓冲区中的消息批量发送到Kafka集群。重试机制:如果消息发送失败(如网络问题、Broker故障等),生产者会根据retries参数的设置进行重试。性能调优:通过调整batch.size、linger.ms等参数,可以优化生产者的性能和吞吐量。综上所述,Kafka生产者流程是一个涉及多个步骤和组件的复杂过程。通过合理配置和优化生产者的行为,可以实现高效、可靠的消息发送。
-
MRS是安全模式,kakfa集群把Ranger鉴权停了也连不上,测试报未知错误,但是kafka在客户端中是可以正常使用的。
-
Kafka作为一种分布式消息队列系统,在大数据领域和实时数据处理中扮演着重要的角色。随着Kafka的广泛应用,用户对其功能的需求也在不断增加。延时操作作为其中之一,为用户提供了更多的灵活性和实用性。本文将介绍Kafka中延时操作的相关内容,包括其背后的原理、实现方式以及应用场景。Kafka延时操作的原理Kafka延时操作的实现原理主要基于两个核心组件:Producer和Consumer。在传统的消息队列系统中,消息被发送后立即可被消费者接收,而Kafka的延时操作则在此基础上进行了扩展,允许用户在发送消息时设置延时参数,使得消息在一定时间后才能被消费者消费。具体来说,Kafka中的延时操作主要通过以下步骤实现:消息发送:Producer将消息发送到Kafka集群中的Topic。延时设置:在消息发送的同时,Producer可以设置延时参数,指定消息在多长时间后可被消费者消费。消息存储:Kafka将延时消息存储在Topic的分区中,但并不立即将其发送给消费者。定时器管理:Kafka内部维护了一个定时器管理器,定期检查消息的延时时间是否到期。消息推送:当消息的延时时间到期后,Kafka将消息推送给对应的消费者进行消费。通过以上步骤,Kafka实现了对延时消息的有效管理和推送,确保消息能够在指定的时间点被消费者接收。Kafka延时操作的实现方式Kafka延时操作的实现方式通常依赖于两种机制:基于时间戳的延时和基于特殊Topic的延时。基于时间戳的延时:这种方式是通过设置消息的时间戳来实现延时操作。Producer在发送消息时,可以为消息设置一个未来的时间戳,指定消息在该时间点之后才能被消费者消费。Kafka会根据消息的时间戳进行延时推送,直到时间点到达后才将消息发送给消费者。基于特殊Topic的延时:另一种方式是通过创建专门的延时Topic来实现延时操作。用户可以将需要延时的消息发送到延时Topic中,然后设置一个定时任务来定期检查延时Topic中的消息,并将到期的消息转发到目标Topic供消费者消费。这两种方式各有优劣,用户可以根据具体需求选择合适的实现方式。Kafka延时操作的应用场景Kafka延时操作在实际应用中具有广泛的应用场景,主要包括以下几个方面:消息调度:延时操作可以用于实现消息的定时发送,例如定时提醒、定时任务等。用户可以将需要延时发送的消息发送到Kafka中,然后设置延时参数,使得消息在指定时间点被发送给消费者。重试机制:延时操作还可以用于实现消息的重试机制。当某个消息发送失败时,可以将该消息发送到延时Topic中,并设置一定的延时时间,等待一段时间后再次尝试发送。这样可以有效地降低消息发送失败的概率,提高系统的可靠性。流量控制:延时操作还可以用于实现流量控制,避免系统因突发大量消息而崩溃。通过设置延时参数,可以在系统负载过高时将部分消息延时发送,从而平滑处理系统压力。业务流程控制:延时操作还可以用于实现复杂的业务流程控制,例如订单超时处理、用户活动提醒等。通过设置延时参数,可以在特定的时间点触发相应的业务流程,从而实现自动化的业务处理。
-
寻农业物联网平台嵌入式开发更新
-
最近遇到个问题,数据上游推送到carbon的数据是实时的,大概5分钟一批。但是carbon数据库不知道怎么才能利用检测工具实时抽取数据到kafka中。有大佬帮忙给个建议吗?
-
生产者代码:报错情况:
-
springboot配置多kafkakafka,说起这个玩意,大家应该都不陌生,不知道啥是kafka的直接去搜索就像MySQL的使用一样,我们在用kafka的时候,也会碰到需要使用多个kafka的情况,比如我从kafkaA的一个topic里消费出数据,进行一次处理,然后我要写入到kafkaB的topic里从网上找了很多配置多kafka的教程,但是都不大好使,后来还是找到了,加上我自己改了点点东西,贴出来和大家分享一下~首先是pom文件,kafka的依赖1234<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId></dependency>配置文件spring: kafka: listener: concurrency: 10 one: # kafka地址 bootstrap-servers: 192.168.217.117:9092 producer: # 生产者每次发送消息的时间间隔(毫秒) linger-ms: 5000 # 单条消息最大值(字节) max-request-size: 16384 # 每次批量发送消息的数量 batch-size: 16384 # 缓存区大小 buffer-memory: 33554432 # 队列 topic: test consumer: # 是否自动提交 enable-auto-commit: true # 队列 topic: topic # group ID group-id: consumer two: # kafka地址 bootstrap-servers: 192.168.217.160:9092 producer: # 生产者每次发送消息的时间间隔(毫秒) linger-ms: 100 # 单条消息最大值(单位为字节,且大小指的是经过序列化后的大小) max-request-size: 16384 # 每次批量发送消息的数量 batch-size: 16384 # 缓存区大小 buffer-memory: 33554432 # 队列 topic: test consumer: # 是否自动提交 enable-auto-commit: true # 队列 topic: test # group ID group-id: consumer有了配置文件当然就要读取配置文件了先读取第一个kafka@EnableKafka @Configuration public class KafkaConfigOne { @Value("${spring.kafka.one.bootstrap-servers}") private String bootstrapServers; @Value("${spring.kafka.one.consumer.enable-auto-commit}") private boolean enableAutoCommit; @Value("${spring.kafka.one.consumer.group-id}") private String groupId; @Value("${spring.kafka.one.producer.linger-ms}") private Integer lingerMs; @Value("${spring.kafka.one.producer.max-request-size}") private Integer maxRequestSize; @Value("${spring.kafka.one.producer.batch-size}") private Integer batchSize; @Value("${spring.kafka.one.producer.buffer-memory}") private Integer bufferMemory; @Bean public KafkaTemplate<String, String> kafkaOneTemplate() { return new KafkaTemplate<>(producerFactory()); } @Bean KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaOneContainerFactory() { ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setConcurrency(3); factory.getContainerProperties().setPollTimeout(3000); return factory; } private ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } public ConsumerFactory<Integer, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); } private Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.LINGER_MS_CONFIG,lingerMs); props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, maxRequestSize); props.put(ProducerConfig.BATCH_SIZE_CONFIG,batchSize); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG,bufferMemory); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; } private Map<String, Object> consumerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit); props.put(ConsumerConfig.GROUP_ID_CONFIG,groupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return props; } } 接着读取第二个kafka@EnableKafka @Configuration public class KafkaConfigTwo { @Value("${spring.kafka.two.bootstrap-servers}") private String bootstrapServers; @Value("${spring.kafka.two.consumer.enable-auto-commit}") private boolean enableAutoCommit; @Value("${spring.kafka.two.consumer.group-id}") private String groupId; @Value("${spring.kafka.two.producer.linger-ms}") private Integer lingerMs; @Value("${spring.kafka.two.producer.max-request-size}") private Integer maxRequestSize; @Value("${spring.kafka.two.producer.batch-size}") private Integer batchSize; @Value("${spring.kafka.two.producer.buffer-memory}") private Integer bufferMemory; @Bean public KafkaTemplate<String, String> kafkaTwoTemplate() { return new KafkaTemplate<>(producerFactory()); } @Bean KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaTwoContainerFactory() { ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setConcurrency(3); factory.getContainerProperties().setPollTimeout(3000); return factory; } private ProducerFactory<String, String> producerFactory() { return new DefaultKafkaProducerFactory<>(producerConfigs()); } public ConsumerFactory<Integer, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); } private Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.LINGER_MS_CONFIG,lingerMs); props.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, maxRequestSize); props.put(ProducerConfig.BATCH_SIZE_CONFIG,batchSize); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG,bufferMemory); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; } private Map<String, Object> consumerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, enableAutoCommit); props.put(ConsumerConfig.GROUP_ID_CONFIG,groupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return props; } } 使用:12345 @Autowired @Qualifier("kafkaTwoTemplate") private KafkaTemplate kafkaTwoTemplate; kafkaTwoTemplate.send("topic", "message");直接注入使用就可以,但是,在这儿注意,因为配置了多个kafka,所以需要进行区分,此处我使用@Autowired和@Qualifier连用,大家也可以使用@DataSource,这个就是多数据源的注解而已,无所谓,按照个人习惯进行使用就OK当然,还有消费者1234@KafkaListener(topics = {"#{'${spring.kafka.two.consumer.topic}'}"}, containerFactory = "kafkaTwoContainerFactory") public void listenerTwo (String data) { System.out.println(data); }
-
FusionInsight HD 6513 在线升级 FusionInsight HD 6517版本 需要多长时间?怎么评估的?
-
FusionInsight HD 6513升级 FusionInsight HD 6517版本,是否支持部分组件在线升级,其他组件离线升级?
-
kafka消息格式,参考cid:link_0带有table、sql、old等关键字,flinksql建表时sql校验不通过,这种情况要如何处理
推荐直播
-
华为云码道Skill实战与极速交付,智能开发全链路实战2026/07/22 周三 19:00-21:00
王一男-华为云码道产品规划专家;李炎-华为云码道产品专家;姜浩-华为云HCDG核心组成员
直播深度解读华为云码道6月产品新特性,从Skill市场安装专家技能,带你零距离体验从需求,开发,审查,重构全链路闭环的开发过程。从零构建并交付一个完整项目,让您体验从代码提交到服务上线的“极速”之旅。
回顾中 -
聚开发者之力,创具身新未来2026/07/23 周四 15:00-17:00
张豪杰/程文/王军/刘新春/黄钦开 /张晓天
本次华为云具身智能开发平台CloudRobo培训面向具身智能开发者,带您全流程体验机器人本体R2C小时级接入、环境重建与轨迹生成仿真数据生产、PB级数据管理、数据评测、模型训推、强化学习和Benchmark一键评测等功能,并体验业界主流具身模型应用。
回顾中
热门标签