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

如何从Java模型对象中提取指定字段列表转JSON发送至Kafka

实现方案

有两种常用的可落地实现方式,都可以完全满足动态筛选字段序列化的需求:


方案1:Jackson动态字段过滤(推荐)

通过Jackson自带的过滤注解配置序列化规则,无需额外对象转换,性能损耗更低。

  1. 给Transaction实体类添加序列化过滤注解:
import com.fasterxml.jackson.annotation.JsonFilter;

@JsonFilter("dynamicTransactionFilter")
public class Transaction{
    String transctionId;
    String accountId;
    String transName;
    String accountName;
    // 补充所有字段的getter、setter方法,Jackson序列化需要读取字段值
}
  1. 构造适配动态字段的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);
  1. 遍历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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 03:36:00