Spring Boot集成Pulsar调用Message.getValue()疑似死锁求助
Spring Boot Pulsar 消费者无法访问消息Payload(疑似死锁)
问题描述
- Spring Boot连接Pulsar消息中间件后,消费者能收到消息,但调用
message.getPayload()或message.getValue()时无任何输出,疑似出现死锁 - 初始阶段无法接收消息,调整
ackMode = AckMode.RECORD后解决了重复投递问题,但Payload访问异常仍存在 - 已尝试排查手段:谷歌搜索同类问题、插入调试语句、切换Spring Boot与Pulsar的Message API实现、开启DEBUG模式日志(未发现异常)
- 开发环境:IntelliJ 2023.1(社区版)
完整代码
package dan.teachingagency.covernotices; import com.fasterxml.jackson.databind.ObjectMapper; import dan.teachingagency.Exception.MapperException; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.pulsar.annotation.EnablePulsar; import org.springframework.pulsar.annotation.PulsarListener; import org.springframework.pulsar.listener.AckMode; import org.springframework.stereotype.Service; @Service @SpringBootApplication @EnablePulsar public class CovernoticesApplication { public static void main(String[] args) { SpringApplication.run(CovernoticesApplication.class, args); } // 调试用计数器 private int messageCount = 0; @Autowired public TeacherRepository teacherRepository; ObjectMapper mapper = new ObjectMapper(); @PulsarListener( subscriptionName = "TeacherAdminSub", topics = "persistent://school1/admin_topics/userdetails", schemaType = SchemaType.JSON, ackMode = AckMode.RECORD) void listen(Message<String> message) { // 初始阶段消息无法接收,调整ackMode后解决重复投递,但Payload访问异常 messageCount++; System.out.println("["+new java.util.Date(System.currentTimeMillis())+"] Message ["+messageCount+"]received."); // 以下代码取消注释后无输出,疑似调用getPayload/getValue时死锁 // System.out.println("["+new java.util.Date(System.currentTimeMillis())+"] Message received ["+this.messageCount+"]: ["+message.getPayload()+"]."); // System.out.println("Message value: ["+message.getValue()+"]"); /*Map<String,String> messageProperties = message.getProperties(); String messageType = messageProperties.get("TYPE"); System.out.println("Message type: ["+messageType+"]"); if("TEACHER".equals(messageType)) { System.out.println("Saving message now. ["+message.getValue()+"]"); StaffMember teacher = fromJson(message.getValue(), StaffMember.class); System.out.println("Teacher=["+teacher+"]"); teacherRepository.save(teacher); System.out.println("Saved message:[" + message.getValue() + "]"); } else { System.out.println("Received rogue message ["+message.getValue()+"]"); }*/ } /** * JSON转对象 * @param json JSON字符串 * @param clazz 目标类 * @param <T> 泛型类型 * @return 目标类实例 */ private <T> T fromJson(String json, Class<T> clazz) { try { return mapper.readValue(json, clazz); } catch (Exception e) { throw new MapperException(e.getMessage()); } } }
排查与解决方案
1. 修复Schema类型不匹配问题
你在@PulsarListener中指定了schemaType = SchemaType.JSON,但方法参数用的是Message<String>,这会导致Schema解析冲突。Pulsar的JSON Schema需要对应具体实体类,而非String类型,可尝试:
- 直接将方法参数改为
StaffMember实体类,让Spring Pulsar自动完成反序列化:
@PulsarListener( subscriptionName = "TeacherAdminSub", topics = "persistent://school1/admin_topics/userdetails", schemaType = SchemaType.JSON, ackMode = AckMode.RECORD) void listen(StaffMember staffMember) { messageCount++; System.out.println("["+new java.util.Date(System.currentTimeMillis())+"] Message ["+messageCount+"]received."); System.out.println("Received staff member: " + staffMember.toString()); // 后续业务逻辑 }
- 确保
StaffMember类有无参构造函数,且字段与消息JSON结构完全匹配。
2. 避免混用Pulsar原生Message与Spring封装类
Spring Pulsar对Pulsar原生Message做了封装,直接调用原生API可能引发线程阻塞。如果必须使用Message对象,改用Spring Pulsar提供的org.springframework.pulsar.core.PulsarMessage:
import org.springframework.pulsar.core.PulsarMessage; // ... void listen(PulsarMessage<String> message) { String payload = message.getPayload(); System.out.println("Payload: " + payload); }
3. 添加异常捕获排查潜在问题
在调用Payload相关方法时添加try-catch块,打印完整堆栈信息,确认是否有未捕获的异常:
void listen(Message<String> message) { messageCount++; System.out.println("["+new java.util.Date(System.currentTimeMillis())+"] Message ["+messageCount+"]received."); try { String payload = message.getPayload(); System.out.println("Payload: " + payload); } catch (Exception e) { e.printStackTrace(); // 打印详细异常堆栈 } }
4. 调整消费者线程池配置
Spring Pulsar默认线程池配置可能导致阻塞,可在application.properties中调整:
spring.pulsar.consumer.threads=5 spring.pulsar.consumer.listener-thread-pool-core-size=5 spring.pulsar.consumer.listener-thread-pool-max-size=10
5. 确认版本兼容性
确保spring-boot-starter-pulsar与Pulsar客户端版本匹配,避免因版本不兼容导致底层API异常。例如Spring Boot 3.2.x对应:
<!-- pom.xml --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-pulsar</artifactId> <version>3.2.0</version> </dependency>
内容的提问来源于stack exchange,提问作者Dan
相关产品推荐
相关产品推荐

