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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:27:05