能否用单个Apache Kafka流应用处理同一主题下的多JSON格式?
问题
同一个Kafka主题中存在两种JSON格式的数据:
- 第一种JSON(注:原JSON存在语法错误,
values应为对象而非数组,已修正):
{ "id": 123, "name": "name", "values": { "one": 1, "two": 2, "three": 3, "four": 4 } }
- 第二种JSON:
{ "id": 8, "title": "Microsoft Surface Laptop 4", "description": "Style and speed. Stand out on ...", "price": 1499, "discountPercentage": 10.23, "rating": 4.43, "stock": 68, "brand": "Microsoft Surface", "category": "laptops", "thumbnail": "https://dummyjson.com/image/i/products/8/thumbnail.jpg", "images": [ "https://dummyjson.com/image/i/products/8/1.jpg", "https://dummyjson.com/image/i/products/8/2.jpg", "https://dummyjson.com/image/i/products/8/3.jpg", "https://dummyjson.com/image/i/products/8/4.jpg", "https://dummyjson.com/image/i/products/8/thumbnail.jpg" ] }
希望通过单个Java流应用对这两种JSON分别执行转换操作,再发送至不同主题,是否可行?具体该如何实现?
解决方案
完全可以通过单个Kafka Streams应用实现,核心思路是先按特征区分消息类型,再分支处理,具体步骤如下:
1. 定义数据模型类
为两种JSON格式创建对应的Java实体类,配合Jackson完成序列化/反序列化:
- 第一种数据模型(
TypeA.java):
public class TypeA { private int id; private String name; private Map<String, Integer> values; // 必须提供无参构造器、Getters和Setters }
- 第二种数据模型(
Product.java):
public class Product { private int id; private String title; private String description; private double price; private double discountPercentage; private double rating; private int stock; private String brand; private String category; private String thumbnail; private List<String> images; // 必须提供无参构造器、Getters和Setters }
2. 初始化Kafka Streams配置
配置应用基础参数,指定序列化器:
Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "multi-type-stream-app"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
3. 读取输入主题并分支消息
使用branch()方法根据消息特征区分类型,这里通过判断是否包含专属字段来区分:
StreamsBuilder builder = new StreamsBuilder(); KStream<String, String> inputStream = builder.stream("input-topic"); // 定义分支条件:通过专属字段判断消息类型 Predicate<String, String> isTypeA = (key, value) -> value.contains("\"name\""); Predicate<String, String> isProduct = (key, value) -> value.contains("\"title\""); // 分支为两个独立流 KStream<String, String>[] branches = inputStream.branch(isTypeA, isProduct); KStream<String, String> typeAStream = branches[0]; KStream<String, String> productStream = branches[1];
若存在字段值干扰(比如
description中包含name),可改用JSON Path做精准判断,例如JsonPath.read(value, "$.name") != null。
4. 分别执行转换并发送到对应主题
对每个分支流独立处理转换逻辑,再输出到目标主题:
- 处理TypeA消息:
ObjectMapper objectMapper = new ObjectMapper(); KStream<String, String> transformedTypeA = typeAStream .mapValues(value -> { try { TypeA typeA = objectMapper.readValue(value, TypeA.class); // 自定义转换逻辑:比如给name添加前缀 typeA.setName("processed-" + typeA.getName()); return objectMapper.writeValueAsString(typeA); } catch (IOException e) { // 反序列化失败的消息可返回null,后续过滤丢弃 return null; } }) .filter((key, value) -> value != null); // 过滤转换失败的消息 // 发送到TypeA专属输出主题 transformedTypeA.to("typeA-output-topic");
- 处理Product消息:
KStream<String, String> transformedProduct = productStream .mapValues(value -> { try { Product product = objectMapper.readValue(value, Product.class); // 自定义转换逻辑:比如计算折后价 product.setPrice(product.getPrice() * (1 - product.getDiscountPercentage() / 100)); return objectMapper.writeValueAsString(product); } catch (IOException e) { return null; } }) .filter((key, value) -> value != null); // 发送到Product专属输出主题 transformedProduct.to("product-output-topic");
5. 启动流应用
KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); // 注册关闭钩子,保证应用优雅停止 Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
可选优化
- 配置死信主题,将反序列化失败的消息转发到专门主题,便于问题排查;
- 自定义
Serde类,避免在mapValues中重复创建ObjectMapper; - 使用Schema Registry管理JSON Schema,实现更可靠的类型判断。
内容的提问来源于stack exchange,提问作者Geek
相关产品推荐
相关产品推荐

