如何从Java模型对象中提取指定字段列表转JSON发送至Kafka
实现方案
有两种常用的可落地实现方式,都可以完全满足动态筛选字段序列化的需求:
方案1:Jackson动态字段过滤(推荐)
通过Jackson自带的过滤注解配置序列化规则,无需额外对象转换,性能损耗更低。
- 给Transaction实体类添加序列化过滤注解:
import com.fasterxml.jackson.annotation.JsonFilter; @JsonFilter("dynamicTransactionFilter") public class Transaction{ String transctionId; String accountId; String transName; String accountName; // 补充所有字段的getter、setter方法,Jackson序列化需要读取字段值 }
- 构造适配动态字段的ObjectMapper:
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ser.impl.SimpleBeanPropertyFilter; import com.fasterxml.jackson.databind.ser.impl.SimpleFilterProvider; import java.util.Arrays; import java.util.List; // 客户端传入的要保留的字段列表 List<String> columnNames = Arrays.asList("transctionId","accountName"); // 配置过滤规则:仅保留columnNames中的字段 SimpleBeanPropertyFilter fieldFilter = SimpleBeanPropertyFilter.filterOutAllExcept(columnNames); SimpleFilterProvider filterProvider = new SimpleFilterProvider() .addFilter("dynamicTransactionFilter", fieldFilter); ObjectMapper objectMapper = new ObjectMapper(); objectMapper.setFilterProvider(filterProvider);
- 遍历Stream转换为指定格式JSON发送到Kafka:
import org.springframework.kafka.core.KafkaTemplate; import com.fasterxml.jackson.core.JsonProcessingException; // 假设已经注入KafkaTemplate private KafkaTemplate<String, String> kafkaTemplate; transctions.forEach(transaction -> { try { String targetJson = objectMapper.writeValueAsString(transaction); // 发送到指定Kafka主题 kafkaTemplate.send("your_topic_name", targetJson); } catch (JsonProcessingException e) { // 自行补充序列化异常处理逻辑 e.printStackTrace(); } });
方案2:转Map后筛选字段
无需修改实体类添加注解,适配性更强,适合不能修改实体类定义的场景。
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.Map; import java.util.function.Function; import java.util.stream.Collectors; ObjectMapper objectMapper = new ObjectMapper(); List<String> columnNames = Arrays.asList("transctionId","accountName"); transctions.forEach(transaction -> { try { // 先把Transaction对象转为包含所有字段的Map Map<String, Object> allFieldMap = objectMapper.convertValue(transaction, new TypeReference<Map<String, Object>>() {}); // 筛选出仅包含指定字段的Map Map<String, Object> filteredMap = columnNames.stream() .filter(allFieldMap::containsKey) .collect(Collectors.toMap(Function.identity(), allFieldMap::get)); // 转JSON后发送Kafka String targetJson = objectMapper.writeValueAsString(filteredMap); kafkaTemplate.send("your_topic_name", targetJson); } catch (Exception e) { // 异常处理逻辑 e.printStackTrace(); } });
注意事项
- 两种方案都可以动态适配
columnNames的取值,客户端传入不同的字段列表不需要修改核心逻辑 - 实体类必须补充对应字段的public getter方法,否则Jackson无法正常读取字段值
- 如果需要自定义序列化后的字段名,可以在对应实体字段上添加
@JsonProperty("自定义别名")注解配置
内容的提问来源于stack exchange,提问作者Debo
相关产品推荐
相关产品推荐

