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

基于Kafka ConsumerRecord构建分主题CSV文件的实现疑问

问题背景

我使用KafkaConsumer读取多个Kafka主题的消息,采用自定义Deserializer,核心方法如下:

public Object deserialize(String s, byte[] bytes) {
    if (!isEmpty(bytes)) {
      try {
        return objectMapper.readValue(bytes, new TypeReference<Map<String, Object>>() {
        });
      } catch (Exception e) {
        LOGGER.error("Error while deserializing bytes[] {} to json object ", new String(bytes), e);
      }
    }

    return null;
}

该反序列化过程无报错,获取到的ConsumerRecords中,ConsumerRecord的value字段为Map结构。

需求
  • 利用主题中的消息生成CSV文件
  • 每个主题单独生成对应文件

初步思路:按主题拆分消息,构建Map<String, List<Map<String, Object>>>,键为主题名称,值为对应ConsumerRecord的value集合。

当前疑问
  1. 如何将ConsumerRecord的value以Map形式提取?目前实现的代码只能得到消息列表:
List customerMigrationKafkaMessages =
consumerRecords.stream()
.flatMap(c -> StreamSupport.stream(c.spliterator(), false).map(ConsumerRecord::value))
.collect(Collectors.toList());
  1. 如何遍历ConsumerRecords,构建上述以主题为键、对应消息Map列表为值的Map?
  2. 为何需要将消息转为HashMap?因为Map的键集可作为CSV表头,值集可作为CSV行数据。
补充尝试问题
  • 曾尝试使用Spring的JsonDeserializer(配置kafka.value.deserializer=org.springframework.kafka.support.serializer.JsonDeserializer),但报错:
Caused by: java.lang.IllegalStateException: No type information in headers and no default type provided
  • 已知各主题消息的JSON结构,但因主题较多,若为每个主题创建POJO,后续仍需转为字符串写入文件,且需在反序列化前维护主题与对应类的映射,示例代码如下:
Map myMap = new HashMap<>();
myMap.put("TopicA", MyClassA.class);
myMap.put("TopicB", MyClassB.class);
public Object deserialize(String topic, String s, byte[] bytes) {
  if (!isEmpty(bytes)) {
    try {
      return objectMapper.readValue(bytes, myMap.get(topic));
    } catch (Exception e) {
      LOGGER.error("Error while deserializing bytes[] {} to json object ", new String(bytes), e);
    }
  }
  return null;
}
  • 需求是每个主题对应一个CSV文件,而非每条消息对应一个文件,处理4个主题则生成4个文件。
实现思路

1. 提取ConsumerRecord的value为Map

当前代码只需指定泛型List<Map<String, Object>>,并对value()做类型转换即可:

List<Map<String, Object>> customerMigrationKafkaMessages =
consumerRecords.stream()
.flatMap(c -> StreamSupport.stream(c.spliterator(), false))
.map(record -> (Map<String, Object>) record.value())
.collect(Collectors.toList());

2. 构建主题到消息列表的Map

使用Stream的Collectors.groupingBy直接按主题分组,同时提取每个record的value为Map,无需额外遍历:

Map<String, List<Map<String, Object>>> topicToMessagesMap =
consumerRecords.stream()
.flatMap(c -> StreamSupport.stream(c.spliterator(), false))
.collect(Collectors.groupingBy(
    ConsumerRecord::topic,
    Collectors.mapping(record -> (Map<String, Object>) record.value(), Collectors.toList())
));

3. 转为HashMap的必要性

Map的键集合可直接作为CSV表头(需注意键的顺序,若需固定顺序可提前定义或排序);每个Map的value集合对应CSV的一行数据,按表头顺序取值即可生成行,这种方式无需为每个主题定义POJO,灵活性更高。

4. 生成CSV文件

针对topicToMessagesMap中的每个键值对:

  • 以主题名作为CSV文件名(例如${topicName}.csv)
  • 提取表头:若该主题所有消息结构一致,取第一条消息的Map键集合;若存在结构不一致,合并所有消息的键作为表头,缺失字段填空值
  • 遍历消息列表,按表头顺序取出对应值,拼接成CSV行
  • 使用专业CSV工具类(如Apache Commons CSV、OpenCSV)处理文件写入,避免手动拼接出现格式问题(如字段包含逗号、引号)

5. Spring JsonDeserializer报错解决

若要使用Spring的JsonDeserializer,只需配置默认反序列化类型:

spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.value.default.type=java.util.Map

这样默认将消息反序列化为Map,和自定义Deserializer效果一致,无需额外维护主题与类的映射。

6. POJO方案的取舍

主题数量较多时,不建议为每个主题创建POJO:维护成本高,且生成CSV时仍需将POJO转为键值对提取字段,Map结构更直接高效。

内容的提问来源于stack exchange,提问作者Abhinash Jha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 06:43:15