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

Spring Boot集成Pulsar Client发送JsonArray批量数据报错排查

问题分析

你遇到的JsonArray.getAsInt错误,核心原因有两个:

  1. JSON库类冲突:代码中同时导入了net.sf.json.JSONArray和org.apache.pulsar.shade.com.google.gson.JsonArray两个不同的JSON数组实现,类型混淆导致解析异常。
  2. Pulsar批量发送的误解:你当前是把包含多个用户的JsonArray作为单条消息发送,而Pulsar的生产者批量是指自动将多条独立消息打包发送,并非单条消息内的数组结构。
解决方案

方案1:使用POJO类(推荐)

定义用户POJO,让Pulsar的JSONSchema直接处理Java对象,避免JSON库类型问题:

// 定义User实体类
public class User {
    private Integer userId;
    private String firstName;

    // 必须提供无参构造器,JSONSchema需要
    public User() {}

    public User(Integer userId, String firstName) {
        this.userId = userId;
        this.firstName = firstName;
    }

    // getter和setter
    public Integer getUserId() { return userId; }
    public void setUserId(Integer userId) { this.userId = userId; }
    public String getFirstName() { return firstName; }
    public void setFirstName(String firstName) { this.firstName = firstName; }
}

修改生产者代码,批量发送独立的User消息:

@Component
@RequiredArgsConstructor
@Slf4j
public class PulsarProducer {

  private static final String TOPIC_NAME = "Json_Test";
  private final PulsarClient client;

  @Bean(name = "producer")
  public void producer() throws PulsarClientException {
    // 批量配置:当攒够2条消息或等待60秒时发送批次
    Producer<User> producer = client.newProducer(JSONSchema.of(User.class))
        .topic(TOPIC_NAME)
        .batchingMaxPublishDelay(60, TimeUnit.SECONDS)
        .batchingMaxMessages(2)
        .enableBatching(true)
        .compressionType(CompressionType.LZ4)
        .create();

    // 构造用户列表,批量发送每条User消息
    List<User> users = Arrays.asList(
        new User(1, "AAAAA"),
        new User(2, "BBBB"),
        new User(3, "CCCCC"),
        new User(4, "DDDDD"),
        new User(5, "EEEEE")
    );

    try {
      for (User user : users) {
        producer.send(user);
        log.info("Sent user: {}", user.getUserId());
      }
    } catch (Exception e) {
      log.error("Error sending messages", e);
    } finally {
      producer.close();
    }
  }
}

方案2:改用String类型Schema发送JSON字符串

如果坚持使用JSON数组格式,直接发送字符串形式的JSON,避免JSON库类型冲突:

@Component
@RequiredArgsConstructor
@Slf4j
public class PulsarProducer {

  private static final String TOPIC_NAME = "Json_Test";
  private final PulsarClient client;

  @Bean(name = "producer")
  public void producer() throws PulsarClientException {
    Producer<String> producer = client.newProducer(Schema.STRING)
        .topic(TOPIC_NAME)
        .batchingMaxPublishDelay(60, TimeUnit.SECONDS)
        .batchingMaxMessages(2)
        .enableBatching(true)
        .compressionType(CompressionType.LZ4)
        .create();

    String data = "[{'userId': 1,'firstName': 'AAAAA'},{'userId': 2,'firstName': 'BBBB'},{'userId': 3,'firstName': 'CCCCC'},{'userId': 4,'firstName': 'DDDDD'},{'userId': 5,'firstName': 'EEEEE'}]";
    
    try {
      producer.send(data);
    } catch (Exception e) {
      log.error("Error sending message", e);
    } finally {
      producer.close();
    }
  }
}

方案3:修复JSON库冲突(不推荐)

如果必须使用Gson的JsonArray,删除net.sf.json.JSONArray的导入,确保全程使用Pulsar shading的Gson类:

// 删除此行导入
// import net.sf.json.JSONArray;

同时验证消费者是否也使用相同的Gson类解析,否则仍会出现类型不匹配问题。

关键注意点
  • Pulsar的生产者批量是多条消息的打包发送,不是单条消息内的数组结构;如果需要发送批量数据作为单条消息,直接发送字符串或POJO列表即可。
  • 避免同时导入多个JSON库,优先使用Pulsar依赖中已包含的Gson(org.apache.pulsar.shade.com.google.gson),防止类加载冲突。

内容的提问来源于stack exchange,提问作者Hyungin Cho

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 12:15:33