Проблема:
- В режиме одной записи (по умолчанию) сопоставление JSON работает отлично, используя Spring.json.type.mapping, и экземпляр абстрактного класса правильно создается в его конкретных реализациях.
- В пакетном режиме потребитель получает List (или пустой список, если я принудительно использую тип DTO) вместо List. Кажется, CompositeMessageConverter или JsonDeserializer не запускает полиморфное сопоставление при переносе в список.
Код: Выделить всё
spring:
cloud:
stream:
bindings:
batchConsumer-in-0:
destination: my-topic
group: my-group
consumer:
batch-mode: true
useNativeDecoding: false # I tried both true and false
kafka:
binder:
consumer-properties:
spring.json.type.mapping: "typeA:com.example.ConcreteA,typeB:com.example.ConcreteB"
spring.json.trusted.packages: "com.example"
Код: Выделить всё
@Bean
public Consumer[*]> batchConsumer() {
return payload -> {
// Here, 'payload' contains byte[] or is incorrectly mapped
log.info("Received batch of size: {}", payload.size());
};
}
Код: Выделить всё
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
4.0.0
org.springframework.boot
spring-boot-starter-parent
3.5.12
21
2025.0.1
1.6.3
^
true
kafka-consumer
Код: Выделить всё
org.springframework.boot
spring-boot-starter
de.codecentric
spring-boot-admin-starter-client
3.5.8
org.springframework.boot
spring-boot-starter-web
org.springframework.boot
spring-boot-starter-actuator
com.github.ulisesbocchio
jasypt-spring-boot-starter
3.0.5
org.springframework.kafka
spring-kafka
org.apache.kafka
kafka-streams
org.springframework.cloud
spring-cloud-starter-stream-kafka
org.springframework.cloud
spring-cloud-stream-binder-kafka-streams
org.springframework.cloud
spring-cloud-starter-config
commons-logging
commons-logging
com.melloware
jasypt
1.9.4
org.springframework.cloud
spring-cloud-starter-consul-discovery
org.projectlombok
lombok
org.springframework.boot
spring-boot-starter-test
test
org.springframework.cloud
spring-cloud-starter-loadbalancer
org.springdoc
springdoc-openapi-starter-webmvc-ui
2.8.16
org.springframework.cloud
spring-cloud-dependencies
${spring-cloud.version}
pom
import
- Использование Consumer -> Без изменений.
- Принудительное использование типа контента: application/json в привязке -> Без изменений.
- Настройка Spring.json.value.default.type в абстрактный класс -> По-прежнему возвращает необработанные байты или не создает экземпляр.
- Использование List работает, но я хочу воспользоваться преимуществами автоматической полиморфной десериализации, которая работает в непакетном режиме.
Как Могу ли я заставить Spring Cloud Stream Kafka применить JsonDeserializer с сопоставлением типов к каждому элементу пакета, если целевым типом является абстрактный класс?
Спасибо за вашу помощь!
Мне наконец удалось решить свою проблему, поэтому я делюсь решением на случай, если оно поможет другим, столкнувшимся с аналогичными проблемами с десериализацией Spring Cloud Stream + Kafka + JSON.
Основная причина
Проблема на самом деле заключалась в сочетании дизайна DTO (Джексон) и конфигурации Spring Cloud Stream:
- Отсутствует конструктор по умолчанию в абстрактном классе
- Мой базовый DTO был абстрактным и не имел конструктора без аргументов
/> - При использовании Jackson (и особенно при использовании Lombok) это обязательно.
- Добавление конструктора по умолчанию исправляет ошибки десериализации.
- Мой базовый DTO был абстрактным и не имел конструктора без аргументов
Вот конфигурация, которая работает как для стандартных, так и для пакетных потребителей, включая полиморфную десериализацию с использованием Заголовки __TypeId__.
Код: Выделить всё
spring:
config:
activate:
on-profile: kafka-stream
autoconfigure:
exclude:
- org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration
cloud:
function:
definition: eventConsumer;batchConsumer
stream:
default:
content-type: application/json
consumer:
concurrency: 3
maxAttempts: 1
header-mode: headers
bindings:
eventConsumer-in-0:
destination: digit-message-topic
group: event-group
consumer:
use-native-decoding: false
enable-dlq: true
dlq-name: digit-message-topic.DLT
auto-offset-reset: earliest
batchConsumer-in-0:
destination: digit-message2-topic
group: batch-group
consumer:
batch-mode: true
use-native-decoding: true
header-mode: headers
enable-dlq: true
dlq-name: digit-message2-topic.DLT
auto-offset-reset: earliest
max-poll-records: 20
fetch-min-bytes: 1
kafka:
default:
consumer:
configuration:
key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
value.deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
spring.json.use.type.headers: true
spring.json.trusted.packages: "be.uclouvain.digit.functional.messenger.dto"
spring.json.type.mapping: "typeA:com.example.ConcreteA,typeB:com.example.ConcreteB"
enable-dlq: true
auto-offset-reset: earliest
binder:
auto-create-topics: true
auto-add-partitions: true
configuration:
bootstrap-servers: localhost:9092
logging:
level:
org.springframework.kafka.support.serializer: TRACE
com.fasterxml.jackson: DEBUG
org.springframework.kafka: OFF
- Всегда предоставляйте конструктор без аргументов, особенно для абстрактных DTO
- Используйте JsonDeserializer с:
Код: Выделить всё
spring.json.use.type.headers=true- если вы используете простые имена типов в __TypeId__
Код: Выделить всё
spring.json.type.mapping - Используйте use-native-decoding=true при делегировании десериализации в Kafka
- Пакетный режим требует правильной настройки Kafka (и т. д.)
Код: Выделить всё
max-poll-records - Будьте осторожны с именованием свойств (vs CamelCase)
Код: Выделить всё
kebab-case
После этих изменений:
- Пакетное потребление работает правильно
- DTO правильно десериализованы (включая полиморфизм)
- Больше никаких ошибок JsonDeserializer или пустых пакеты