咨询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
相关产品推荐
相关产品推荐

