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

咨询Spring Cloud KStream消费JSON并生产Avro的配置示例

Spring Cloud Stream Kafka Streams JSON转Avro实现示例

以下是完整的实现代码,包含依赖配置、实体定义、Streams配置和核心处理逻辑:

1. 依赖配置(Maven)

需要引入Spring Cloud Stream Kafka Streams binder、Confluent的JSON和Avro Serde,以及Avro代码生成插件:

<dependencies>
    <!-- Spring Cloud Stream Kafka Streams Binder -->
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-stream-binder-kafka-streams</artifactId>
    </dependency>
    <!-- Confluent JSON Schema Serde -->
    <dependency>
        <groupId>io.confluent</groupId>
        <artifactId>kafka-streams-json-schema-serde</artifactId>
        <version>7.5.0</version> <!-- 与Kafka版本匹配 -->
    </dependency>
    <!-- Confluent Avro Serde -->
    <dependency>
        <groupId>io.confluent</groupId>
        <artifactId>kafka-streams-avro-serde</artifactId>
        <version>7.5.0</version>
    </dependency>
    <!-- Avro API -->
    <dependency>
        <groupId>org.apache.avro</groupId>
        <artifactId>avro</artifactId>
        <version>1.11.3</version>
    </dependency>
</dependencies>

<!-- Avro代码生成插件,用于生成TestAvro类 -->
<build>
    <plugins>
        <plugin>
            <groupId>org.apache.avro</groupId>
            <artifactId>avro-maven-plugin</artifactId>
            <version>1.11.3</version>
            <executions>
                <execution>
                    <phase>generate-sources</phase>
                    <goals>
                        <goal>schema</goal>
                    </goals>
                    <configuration>
                        <sourceDirectory>${project.basedir}/src/main/resources/avro</sourceDirectory>
                        <outputDirectory>${project.build.directory}/generated-sources/avro</outputDirectory>
                    </configuration>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>

2. 实体类定义

TestJson(JSON消息实体)

普通POJO,对应输入主题的JSON结构:

public class TestJson {
    private String id;
    private String name;
    private int age;

    // 必须提供无参构造器,用于JSON反序列化
    public TestJson() {}

    // Getter、Setter方法
    public String getId() { return id; }
    public void setId(String id) { this.id = id; }
    public String getName() { return name; }
    public void setName(String name) { this.name = name; }
    public int getAge() { return age; }
    public void setAge(int age) { this.age = age; }
}

TestAvro(Avro消息实体)

先编写Avro Schema文件src/main/resources/avro/TestAvro.avsc:

{
  "type": "record",
  "name": "TestAvro",
  "namespace": "com.example.demo",
  "fields": [
    {"name": "id", "type": "string"},
    {"name": "fullName", "type": "string"},
    {"name": "userAge", "type": "int"}
  ]
}

通过Avro插件生成的TestAvro类会自动实现SpecificRecord接口,无需手动编写。

3. Kafka Streams配置类

配置自定义的Serde实例,以及Spring Cloud Stream的绑定逻辑:

import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde;
import io.confluent.kafka.streams.serdes.json.KafkaJsonSchemaSerde;
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.util.HashMap;
import java.util.Map;

import static io.confluent.kafka.serializers.AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG;

@Configuration
public class StreamsSerdeConfig {

    // 配置Schema Registry地址(可从配置文件注入)
    private final String schemaRegistryUrl = "http://localhost:8081";

    // JSON Schema Serde配置,用于消费TestJson
    @Bean
    public KafkaJsonSchemaSerde<TestJson> testJsonSerde() {
        KafkaJsonSchemaSerde<TestJson> jsonSerde = new KafkaJsonSchemaSerde<>();
        Map<String, String> configs = new HashMap<>();
        configs.put(SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        configs.put("json.value.type", TestJson.class.getName());
        jsonSerde.configure(configs, false); // false表示是值序列化器(不是键)
        return jsonSerde;
    }

    // Specific Avro Serde配置,用于生产TestAvro
    @Bean
    public SpecificAvroSerde<TestAvro> testAvroSerde() {
        SpecificAvroSerde<TestAvro> avroSerde = new SpecificAvroSerde<>();
        Map<String, String> configs = new HashMap<>();
        configs.put(SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        avroSerde.configure(configs, false);
        return avroSerde;
    }
}

4. 核心流处理逻辑

使用Spring Cloud Stream的函数式编程模型,定义processRecord函数实现JSON到Avro的转换:

import org.apache.kafka.streams.kstream.KStream;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.util.function.Function;

@Configuration
public class StreamProcessingConfig {

    @Bean
    public Function<KStream<String, TestJson>, KStream<String, TestAvro>> processRecord() {
        return inputStream -> inputStream
                // 转换TestJson到TestAvro
                .mapValues(json -> TestAvro.newBuilder()
                        .setId(json.getId())
                        .setFullName(json.getName())
                        .setUserAge(json.getAge())
                        .build())
                // 可选:添加日志或其他处理逻辑
                .peek((key, avro) -> System.out.printf("转换完成:%s -> %s%n", key, avro));
    }
}

5. 配置文件(application.yml)

配置Spring Cloud Stream绑定、Kafka Streams参数和Serde映射:

spring:
  cloud:
    stream:
      function:
        definition: processRecord # 指定要启用的处理函数
      bindings:
        processRecord-in-0: # 输入绑定,对应函数的输入参数
          destination: input-json-topic # 输入主题名称
          content-type: application/json+schema # 使用JSON Schema Serde
          consumer:
            valueSerde: com.example.demo.StreamsSerdeConfig#testJsonSerde # 指定自定义JSON Serde
        processRecord-out-0: # 输出绑定,对应函数的输出结果
          destination: output-avro-topic # 输出主题名称
          content-type: application/vnd.apache.avro+json # 使用Avro Serde
          producer:
            valueSerde: com.example.demo.StreamsSerdeConfig#testAvroSerde # 指定自定义Avro Serde
      kafka:
        streams:
          binder:
            configuration:
              schema.registry.url: http://localhost:8081 # 全局Schema Registry配置
              default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
              # 其他Kafka Streams配置(可选)
              commit.interval.ms: 1000

关键说明

  • Schema Registry:确保Schema Registry服务已启动并运行,JSON和Avro的Schema会自动注册到Registry中(首次生产/消费时)。
  • Serde配置:KafkaJsonSchemaSerde需要指定json.value.type来明确反序列化的目标类;SpecificAvroSerde依赖Avro生成的SpecificRecord类。
  • 函数式模型:Spring Cloud Stream 3.x+推荐使用函数式编程模型,通过function.definition指定要启用的处理函数,自动绑定输入输出主题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 09:53:09