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

能否用单个Apache Kafka流应用处理同一主题下的多JSON格式?

问题

同一个Kafka主题中存在两种JSON格式的数据:

  1. 第一种JSON(注:原JSON存在语法错误,values应为对象而非数组,已修正):
{
  "id": 123,
  "name": "name",
  "values": {
     "one": 1,
     "two": 2,
     "three": 3,
     "four": 4
   }
}
  1. 第二种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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 05:35:23