基于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集合。
当前疑问
- 如何将ConsumerRecord的value以Map形式提取?目前实现的代码只能得到消息列表:
List customerMigrationKafkaMessages = consumerRecords.stream() .flatMap(c -> StreamSupport.stream(c.spliterator(), false).map(ConsumerRecord::value)) .collect(Collectors.toList());
- 如何遍历ConsumerRecords,构建上述以主题为键、对应消息Map列表为值的Map?
- 为何需要将消息转为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
相关产品推荐
相关产品推荐

