Кафка не может десериализовать сложный вложенный классJAVA

Программисты JAVA общаются здесь
Anonymous
Кафка не может десериализовать сложный вложенный класс

Сообщение Anonymous »

Мой вариант использования:
  • Когда добавляется новый адрес, я отправляю объект пользовательского домена в kafka (он включает адрес массива), чтобы уведомить Redis кэширование этих данных.
Но есть проблема: когда я отправляю пользовательские данные со многих адресов, они все равно отправляются в Kafka, и тема все равно получает их, но это перемещает данные в недоставленное сообщение (Успешная публикация недоставленного сообщения: кэш-пользователь-0@88 в кэш-пользователь.DLT-0@60)
Но когда я отправляю пользовательские данные с нулевым или пустым значением, все работает гладко.
Может быть, он не может десериализовать адрес вложенного класса?
Может ли кто-нибудь Помогите мне решить, я уже давно застрял с проблемой?
Заранее спасибо
Мой пользовательский домен:

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

@Getter
@Builder
@AllArgsConstructor
@NoArgsConstructor
public class User {
private Long id;
private String fullName;
private Email email;
private List addresses;
private PhoneNumber phoneNumber;
}
Субъект домена моего адреса:

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

@Getter
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class Address {
private Long id;
private String province;
private String district;
private String ward;
private String homeAddress;
private AddressType type;
private Long userId;
}
Мой потребитель:

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

@KafkaListener(topics = { CACHE_USER_TOPIC }, groupId = KafkaConfiguration.GROUP_ID)
@KafkaHandler(isDefault = true)
public void cacheUser(ConsumerRecord record) {
User inputData = this.objectMapper.convertValue(record.value(),
User.class);
this.userRedisCaching.cache(inputData);
}
Мой вариант использования:

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

@Override
@Transactional("transactionManager")
public UserAddressOutputModel createAddress(Long userId, CreateUserAddressInputModel inputModel) {
Optional matchedUser = this.loadUserPort.loadUser(userId);

if (matchedUser.isEmpty()) {
throw new UserNotFoundException();
}

if (matchedUser.get().getAddresses().size() >= MAX_ADDRESSES) {
throw new ReachTheMaximumAddressesPerUser();
}

Address address = Address.builder()
.province(inputModel.getProvince())
.district(inputModel.getDistrict())
.ward(inputModel.getWard())
.homeAddress(inputModel.getHomeAddress())
.type(new AddressType(inputModel.getType()))
.userId(matchedUser.get().getId())
.build();
Address newAddress = this.createUserAddressPort.createUserAddress(address);
matchedUser.get().addNewAddress(newAddress);
// User user = User.builder().addresses(Collections.emptyList()).build();
this.sendEventToMessageQueuePort.cacheUser(matchedUser.get());
return UserAddressOutputModel.convertFromDomain(newAddress);
}
Мое приложение.свойства

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

spring.kafka.bootstrap-servers=localhost:9092

# Producer Configuration
spring.kafka.producer.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer

# Consumer Configuration
spring.kafka.consumer.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=exam-outline-pj
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.key-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.trusted.packages=*

# kafka transaction
spring.kafka.producer.transaction-id-prefix=tx-
spring.kafka.consumer.enable-auto-commit=false
spring.kafka.consumer.isolation-level=READ_COMMITTED
spring.kafka.listener.ack-mode=RECORD

# kafka debug
logging.level.org.springframework.kafka=DEBUG
Решите проблему и определите ее причину.

Подробнее здесь: https://stackoverflow.com/questions/788 ... sted-class

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