Spring Boot集成Pulsar Client发送JsonArray批量数据报错排查
问题分析
你遇到的JsonArray.getAsInt错误,核心原因有两个:
- JSON库类冲突:代码中同时导入了
net.sf.json.JSONArray和org.apache.pulsar.shade.com.google.gson.JsonArray两个不同的JSON数组实现,类型混淆导致解析异常。 - 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
相关产品推荐
相关产品推荐

