You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 17:39:51