Я разрабатываю приложение Spring Cloud Stream, которое принимает сообщение в формате Avro. После настройки всего я запускаю свое приложение, а затем получаю неизвестную ошибку магического байта. Я думаю, что это потому, что я использую SpecialAvroserde. Любая помощь будет оценена. Вот мой код
application.yml:
spring:
cloud:
function:
definition: employeeSalaryUpdateFlow
stream:
bindings:
employeeSalaryUpdateFlow-in-0:
destination: employee-topic
kafka:
streams:
binder:
applicationId: kafka-streams-app
brokers: localhost:9092
configuration:
schema.registry.url: http://localhost:8081
processing.guarantee: exactly_once
commit.interval.ms: 10000
default:
key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
bindings:
employeeSalaryUpdateFlow-in-0:
consumer:
materializedAs: employeeSalaryUpdateFlow-store
value-serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde
processor.java
@Bean
public Consumer employeeSalaryUpdateFlow() {
SpecificAvroSerde avroSerde = new SpecificAvroSerde();
Map serdeConfig = new HashMap();
serdeConfig.put("schema.registry.url", "http://localhost:8081"); // Set your schema registry URL
avroSerde.configure(serdeConfig, false);
return input -> {
input
.map((k, v) -> {
log.info("Hey");
return KeyValue.pair(v.getId(), v);
})
.toTable()
.groupBy((k, v) -> KeyValue.pair(v.getDepartment(), v), Grouped.with(Serdes.String(), avroSerde))
.aggregate(() -> employeeSalaryRecordBuilder.init(),
(k, v, agg) -> employeeSalaryRecordBuilder.aggregate(v, agg),
(k, v, agg) -> employeeSalaryRecordBuilder.substract(v, agg))
.toStream()
.foreach((k, v) -> log.info(String.format("Department: %s, Total Salary: %f", k, v.getTotalSalary())));
};
}
employee.avro
{
"namespace": "me.model",
"type": "record",
"name": "Employee",
"fields": [
{"name": "id","type": ["null","string"]},
{"name": "name","type": ["null","string"]},
{"name": "department","type": ["null","string"]},
{"name": "salary","type": ["null","int"]}
]
}
Подробнее здесь: https://stackoverflow.com/questions/797 ... ggregation