StreamBridge未发送Kafka确认头问题排查及手动ACK配置咨询
问题现象
使用KafkaTemplate发送消息时能正常携带确认头,但通过StreamBridge发送时无法传递确认头。
业务场景
我们从数据库接收事件,尝试通过StreamBridge将事件发送到消息通道,需要对事件执行手动ACK操作。
相关配置与代码
Application.properties
server.port=8082 spring.cloud.stream.function.definition=sink1;sink2 spring.cloud.stream.function.bindings.sink1-in-0=inbound-events spring.cloud.stream.bindings.inbound-events.group=ama1-channel-group spring.cloud.stream.bindings.inbound-events.destination=squaredNumbers-test4 spring.cloud.stream.bindings.inbound-events.consumer.header-mode=headers spring.cloud.stream.bindings.inbound-events.content-type=application/json spring.cloud.stream.kafka.bindings.inbound-events.consumer.ack-mode=manual spring.cloud.stream.function.bindings.sink2-in-0=inbound-stream spring.cloud.stream.bindings.inbound-stream.group=ama2-channel-group spring.cloud.stream.bindings.inbound-stream.destination=squaredNumbers-test5 spring.cloud.stream.bindings.inbound-stream.consumer.header-mode=headers spring.cloud.stream.bindings.inbound-stream.content-type=application/json spring.cloud.stream.kafka.bindings.inbound-stream.consumer.ack-mode=manual
Service类
@Service public class KafkaConsumer { BindingServiceProperties bindingProperties; StreamBridge streamBridge; @Autowired public KafkaConsumer(final BindingServiceProperties bindingServiceProperties, StreamBridge streamBridge) { this.bindingProperties = bindingServiceProperties; this.streamBridge = streamBridge; } @Bean public Consumer<Message> sink1() { return (message) -> { System.out.println("******************"); System.out.println("At Sink1"); System.out.println("******************"); System.out.println("Received message " + message); streamBridge.send("inbound-stream",MessageBuilder.fromMessage(message)); }; } @Bean public Consumer<Message> sink2() { return (message) -> { System.out.println("******************"); System.out.println("At Sink2"); System.out.println("******************"); System.out.println("Received message " + message); }; } }
主应用类
@SpringBootApplication(scanBasePackages = "demo") public class MultipleFunctionsApplication { public static void main(String[] args) { SpringApplication.run(MultipleFunctionsApplication.class, args); } }
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 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <artifactId>multi-functions-kafka</artifactId> <version>0.0.1-SNAPSHOT</version> <packaging>jar</packaging> <name>multi-functions-kafka</name> <description>Spring Cloud Stream Sample Multiple functions Kafka</description> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.5</version> <relativePath/> <!-- lookup parent from repository --> </parent> <properties> <spring-cloud.version>2021.0.5</spring-cloud.version> </properties> <dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-dependencies</artifactId> <version>${spring-cloud.version}</version> <type>pom</type> <scope>import</scope> </dependency> </dependencies> </dependencyManagement> <dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-stream-binder-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-actuator</artifactId> </dependency> <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> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-stream</artifactId> <scope>test</scope> <classifier>test-binder</classifier> <type>test-jar</type> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> </dependency> <dependency> <groupId>org.apache.avro</groupId> <artifactId>avro</artifactId> <version>1.10.2</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> <repositories> <repository> <id>spring-snapshots</id> <name>Spring Snapshots</name> <url>https://repo.spring.io/libs-snapshot-local</url> <snapshots> <enabled>true</enabled> </snapshots> <releases> <enabled>false</enabled> </releases> </repository> <repository> <id>spring-milestones</id> <name>Spring Milestones</name> <url>https://repo.spring.io/libs-milestone-local</url> <snapshots> <enabled>false</enabled> </snapshots> </repository> </repositories> <pluginRepositories> <pluginRepository> <id>spring-snapshots</id> <name>Spring Snapshots</name> <url>https://repo.spring.io/libs-snapshot-local</url> <snapshots> <enabled>true</enabled> </snapshots> <releases> <enabled>false</enabled> </releases> </pluginRepository> <pluginRepository> <id>spring-milestones</id> <name>Spring Milestones</name> <url>https://repo.spring.io/libs-milestone-local</url> <snapshots> <enabled>false</enabled> </snapshots> </pluginRepository> <pluginRepository> <id>spring-releases</id> <name>Spring Releases</name> <url>https://repo.spring.io/libs-release-local</url> <snapshots> <enabled>false</enabled> </snapshots> </pluginRepository> </pluginRepositories> </project>
使用版本:spring-cloud-dependencies 2021.0.5
问题分析与解决方案
你的配置和代码存在几个关键问题,导致StreamBridge无法传递确认头以及手动ACK失效:
1. StreamBridge目标绑定配置错误
你直接用streamBridge.send("inbound-stream", ...),但inbound-stream是消费者绑定,不是生产者绑定。StreamBridge需要对应生产者的绑定配置,而非复用消费者的destination。
解决:新增StreamBridge专用的生产者绑定,比如:
# 新增StreamBridge使用的生产者绑定 spring.cloud.stream.bindings.outbound-stream.destination=squaredNumbers-test5 spring.cloud.stream.bindings.outbound-stream.producer.header-mode=headers
然后修改StreamBridge发送目标:
streamBridge.send("outbound-stream", MessageBuilder.fromMessage(message).build());
2. 手动ACK逻辑缺失
配置了ack-mode=manual后,必须显式调用Acknowledgment对象的确认方法,否则消息会一直处于未确认状态。
解决:修改Consumer方法,注入Acknowledgment并执行确认:
@Bean public Consumer<Message<?>> sink1() { return (message, ack) -> { System.out.println("******************"); System.out.println("At Sink1"); System.out.println("******************"); System.out.println("Received message " + message); streamBridge.send("outbound-stream", MessageBuilder.fromMessage(message).build()); ack.acknowledge(); // 手动确认消息 }; } @Bean public Consumer<Message<?>> sink2() { return (message, ack) -> { System.out.println("******************"); System.out.println("At Sink2"); System.out.println("******************"); System.out.println("Received message " + message); ack.acknowledge(); // 手动确认消息 }; }
3. 依赖重复问题
pom.xml中存在重复依赖(如spring-boot-starter-test、spring-kafka-test各出现两次),虽然不影响功能,但建议清理重复项,避免潜在冲突。
补充说明
如果需要端到端的ACK(即sink2处理完成后才确认sink1接收到的消息),需要引入消息组或事务机制;如果仅需确认sink1成功接收并转发消息,上述修改即可满足需求。
内容的提问来源于stack exchange,提问作者secret129
相关产品推荐
相关产品推荐

