Spring Cloud Stream Kafka: в пакетном режиме не удается сопоставить JSON с DTO абстрактного класса (возвращает byte[] вмJAVA

Программисты JAVA общаются здесь
Anonymous
Spring Cloud Stream Kafka: в пакетном режиме не удается сопоставить JSON с DTO абстрактного класса (возвращает byte[] вм

Сообщение Anonymous »

Я столкнулся с проблемой с Spring Cloud Stream (Kafka Binder) при использовании пакетного режима: true в сочетании с Абстрактным классом в качестве типа DTO.
Проблема:
  • В режиме одной записи (по умолчанию) сопоставление JSON работает отлично, используя Spring.json.type.mapping, и экземпляр абстрактного класса правильно создается в его конкретных реализациях.
  • В пакетном режиме потребитель получает List (или пустой список, если я принудительно использую тип DTO) вместо List. Кажется, CompositeMessageConverter или JsonDeserializer не запускает полиморфное сопоставление при переносе в список.
Моя конфигурация (application.yml):

Код: Выделить всё

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());
};
}

Мой Pom.xml:

Код: Выделить всё

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) это обязательно.
    • Добавление конструктора по умолчанию исправляет ошибки десериализации.
Рабочая конфигурация

Вот конфигурация, которая работает как для стандартных, так и для пакетных потребителей, включая полиморфную десериализацию с использованием Заголовки __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
  • Код: Выделить всё

    spring.json.type.mapping
    если вы используете простые имена типов в __TypeId__
  • Используйте use-native-decoding=true при делегировании десериализации в Kafka
  • Пакетный режим требует правильной настройки Kafka (

    Код: Выделить всё

    max-poll-records
    и т. д.)
  • Будьте осторожны с именованием свойств (

    Код: Выделить всё

    kebab-case
    vs CamelCase)
Результат
После этих изменений:
  • Пакетное потребление работает правильно
  • DTO правильно десериализованы (включая полиморфизм)
  • Больше никаких ошибок JsonDeserializer или пустых пакеты
Надеюсь, это поможет 👍

Вернуться в «JAVA»