如何通过Spring XML触发Kafka Consumer(无需main方法)
我开发了一个基于Spring XML + Kafka的示例程序,现在想知道如何触发消费者端的逻辑,并且不想通过main方法来启动消费者。相关代码如下:
MainApp.java
public class MainApp { private static final Faker FAKER = Faker.instance().instance(); public static void main(String[] args) throws InterruptedException { ApplicationContext context = new ClassPathXmlApplicationContext("/context.xml"); EmployeeProducer employeeProducer = (EmployeeProducer) context.getBean("employeeProducer"); Employee employee = Employee.builder() .empId(ThreadLocalRandom.current().nextInt(1,100)) .firstName(FAKER.name().firstName()) .lastName(FAKER.name().lastName()) .gender(getGender()) .build(); System.out.println("Employee : "+ employee); employeeProducer.sendMessage("t-employee", employee); Thread.sleep(500_000_000); } private static String getGender(){ int ramdomN = ThreadLocalRandom.current().nextInt(0,1); String sex; if (ramdomN == 0) { sex = "M"; } else { sex = "F"; } return sex; } }
EmployeeProducer.java
public class EmployeeProducer { public void sendMessage(String topicName, Employee employee) { Properties properties = new Properties(); properties.put(ProducerConfig.CLIENT_ID_CONFIG, AppConfigs.applicationID); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, AppConfigs.bootstrapServers); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName()); KafkaProducer<String, Employee> kafkaProducer = new KafkaProducer<>(properties); kafkaProducer.send(new ProducerRecord<>(topicName, employee)); kafkaProducer.close(); } }
EmployeeConsumer.java
public class EmployeeConsumer { public void consumeMessage() { Properties properties = new Properties(); properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, AppConfigs.bootstrapServers); properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class.getName()); properties.put(ConsumerConfig.GROUP_ID_CONFIG, "group1"); properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); properties.put(JsonDeserializer.VALUE_CLASS_NAME_CONFIG, Employee.class); KafkaConsumer<String, Employee> consumer = new KafkaConsumer<>(properties); consumer.subscribe(List.of(AppConfigs.topicName)); while (true) { ConsumerRecords<String, Employee> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, Employee> record : records) { System.out.println("Key: " + record.key() + ", Value:" + record.value()); System.out.println("Partition:" + record.partition() + ",Offset:" + record.offset()); } } } }
context.xml
<?xml version="1.0" encoding="UTF-8"?> <beans xmlns="http://www.springframework.org/schema/beans" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:context="http://www.springframework.org/schema/context" xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-3.0.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-3.0.xsd"> <bean id="employeeProducer" class="com.example.EmployeeProducer" /> <bean id="employeeConsumer" class="com.example.EmployeeConsumer" /> <!-- <bean class="org.springframework.kafka.annotation.KafkaListenerAnnotationBeanPostProcessor" /> <bean class="org.springframework.kafka.config.KafkaListenerEndpointRegistry"/>--> </beans>
pom.xml
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>org.example</groupId> <artifactId>spring-conflient-demo</artifactId> <version>1.0-SNAPSHOT</version> <properties> <maven.compiler.source>17</maven.compiler.source> <maven.compiler.target>17</maven.compiler.target> </properties> <dependencies> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.6.2</version> </dependency> <dependency> <groupId>com.github.javafaker</groupId> <artifactId>javafaker</artifactId> <version>1.0.2</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.datatype</groupId> <artifactId>jackson-datatype-jsr310</artifactId> <version>2.15.4</version> </dependency> <dependency> <groupId>org.apache.logging.log4j</groupId> <artifactId>log4j-slf4j-impl</artifactId> <version>2.19.0</version> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <version>1.18.28</version> </dependency> <dependency> <groupId>commons-lang</groupId> <artifactId>commons-lang</artifactId> <version>2.6</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.13.4</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-core</artifactId> <version>6.0.18</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-web</artifactId> <version>6.0.18</version> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-context</artifactId> <version>6.0.18</version> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>3.1.4</version> </dependency> </dependencies> </project>
利用已引入的spring-kafka依赖,通过Spring原生注解和容器支持自动启动消费者,无需手动编写main方法或调用消费逻辑。具体步骤如下:
1. 改造消费者类,使用@KafkaListener注解
替换原有手动创建KafkaConsumer的逻辑,用Spring Kafka注解监听指定主题,Spring会自动初始化并启动消费者线程:
import org.springframework.kafka.annotation.KafkaListener; public class EmployeeConsumer { @KafkaListener(topics = "t-employee", groupId = "group1") public void consumeMessage(ConsumerRecord<String, Employee> record) { System.out.println("Key: " + record.key() + ", Value:" + record.value()); System.out.println("Partition:" + record.partition() + ",Offset:" + record.offset()); } }
如果只需处理消息内容,可简化参数为Employee employee:
@KafkaListener(topics = "t-employee", groupId = "group1") public void consumeMessage(Employee employee) { System.out.println("Received Employee: " + employee); }
2. 更新Spring XML配置,启用Kafka注解支持
修改context.xml,添加Kafka相关命名空间,配置消费者工厂和监听容器,让Spring识别@KafkaListener注解:
<?xml version="1.0" encoding="UTF-8"?> <beans xmlns="http://www.springframework.org/schema/beans" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:context="http://www.springframework.org/schema/context" xmlns:kafka="http://www.springframework.org/schema/kafka" xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-3.0.xsd http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-3.0.xsd http://www.springframework.org/schema/kafka http://www.springframework.org/schema/kafka/spring-kafka.xsd"> <!-- 扫描组件,让Spring发现@KafkaListener注解 --> <context:component-scan base-package="com.example"/> <!-- Kafka消费者配置 --> <bean id="consumerProps" class="java.util.Properties"> <constructor-arg> <map> <entry key="bootstrap.servers" value="localhost:9092"/> <entry key="key.deserializer" value="org.apache.kafka.common.serialization.StringDeserializer"/> <entry key="value.deserializer" value="org.springframework.kafka.support.serializer.JsonDeserializer"/> <entry key="auto.offset.reset" value="earliest"/> <entry key="spring.json.value.default.type" value="com.example.Employee"/> </map> </constructor-arg> </bean> <!-- 消费者工厂 --> <bean id="consumerFactory" class="org.springframework.kafka.core.DefaultKafkaConsumerFactory"> <constructor-arg ref="consumerProps"/> </bean> <!-- 监听容器工厂,负责创建消费者线程 --> <bean id="kafkaListenerContainerFactory" class="org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory"> <property name="consumerFactory" ref="consumerFactory"/> </bean> <!-- 注册Kafka注解处理器 --> <bean class="org.springframework.kafka.annotation.KafkaListenerAnnotationBeanPostProcessor"/> <!-- 保留生产者Bean --> <bean id="employeeProducer" class="com.example.EmployeeProducer" /> </beans>
建议将Kafka配置参数(如bootstrap.servers)抽入
application.properties,通过<context:property-placeholder location="classpath:application.properties"/>加载,提升配置灵活性。
3. 启动消费者的几种方式
方式一:作为Web应用部署
由于已引入spring-web依赖,可将项目打包为WAR包,部署到Tomcat等Servlet容器中。Spring上下文启动时会自动初始化Kafka消费者并开始监听主题。
方式二:借助Spring Boot启动(推荐)
引入Spring Boot依赖,创建简单启动类,Spring Boot会自动加载上下文并触发Kafka监听:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; @SpringBootApplication public class KafkaConsumerApp { public static void main(String[] args) { SpringApplication.run(KafkaConsumerApp.class, args); } }
需在pom.xml中添加Spring Boot父依赖和启动器:
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.2.5</version> </parent> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> </dependencies>
方式三:Spring生命周期回调启动
若不想用Web容器或Spring Boot,可让消费者类实现InitializingBean接口,在Spring初始化Bean后自动启动消费逻辑:
import org.springframework.beans.factory.InitializingBean; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class EmployeeConsumer implements InitializingBean { private final ExecutorService executor = Executors.newSingleThreadExecutor(); public void consumeMessage() { // 保留原有手动消费逻辑 } @Override public void afterPropertiesSet() throws Exception { // 用线程池启动消费,避免阻塞Spring上下文初始化 executor.submit(this::consumeMessage); } }
4. 生产者优化(可选)
原有生产者每次发送都创建新的KafkaProducer,性能较差。改用spring-kafka的KafkaTemplate优化:
import org.springframework.kafka.core.KafkaTemplate; public class EmployeeProducer { private final KafkaTemplate<String, Employee> kafkaTemplate; // 构造注入KafkaTemplate public EmployeeProducer(KafkaTemplate<String, Employee> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessage(String topicName, Employee employee) { kafkaTemplate.send(topicName, employee); } }
在context.xml中添加生产者配置:
<!-- Kafka生产者配置 --> <bean id="producerProps" class="java.util.Properties"> <constructor-arg> <map> <entry key="bootstrap.servers" value="localhost:9092"/> <entry key="key.serializer" value="org.apache.kafka.common.serialization.StringSerializer"/> <entry key="value.serializer" value="org.springframework.kafka.support.serializer.JsonSerializer"/> <entry key="client.id" value="employee-producer"/> </map> </constructor-arg> </bean> <!-- 生产者工厂 --> <bean id="producerFactory" class="org.springframework.kafka.core.DefaultKafkaProducerFactory"> <constructor-arg ref="producerProps"/> </bean> <!-- KafkaTemplate --> <bean id="kafkaTemplate" class="org.springframework.kafka.core.KafkaTemplate"> <constructor-arg ref="producerFactory"/> </bean>
内容的提问来源于stack exchange,提问作者PAA

