diff --git a/src/main/java/com/ycwl/basic/integration/message/dto/ZtMessage.java b/src/main/java/com/ycwl/basic/integration/message/dto/ZtMessage.java index 2222a343..78f8fc8b 100644 --- a/src/main/java/com/ycwl/basic/integration/message/dto/ZtMessage.java +++ b/src/main/java/com/ycwl/basic/integration/message/dto/ZtMessage.java @@ -15,6 +15,7 @@ import java.util.Map; @AllArgsConstructor @JsonInclude(JsonInclude.Include.NON_NULL) public class ZtMessage { + private String messageId; // unique message identifier private String channelId; // required private String title; // required private String content; // required diff --git a/src/main/java/com/ycwl/basic/integration/message/service/ZtMessageProducerService.java b/src/main/java/com/ycwl/basic/integration/message/service/ZtMessageProducerService.java index 2386a9d5..b7dd2b9c 100644 --- a/src/main/java/com/ycwl/basic/integration/message/service/ZtMessageProducerService.java +++ b/src/main/java/com/ycwl/basic/integration/message/service/ZtMessageProducerService.java @@ -27,18 +27,24 @@ public class ZtMessageProducerService { public void send(ZtMessage msg) { validate(msg); + + // Generate messageId if not present + if (StringUtils.isBlank(msg.getMessageId())) { + msg.setMessageId(java.util.UUID.randomUUID().toString()); + } + String topic = kafkaProps != null && StringUtils.isNotBlank(kafkaProps.getZtMessageTopic()) ? kafkaProps.getZtMessageTopic() : DEFAULT_TOPIC; String key = msg.getChannelId(); String payload = toJson(msg); - log.info("[ZT-MESSAGE] producing to topic={}, key={}, title={}", topic, key, msg.getTitle()); + log.info("[ZT-MESSAGE] producing to topic={}, key={}, messageId={}, title={}", topic, key, msg.getMessageId(), msg.getTitle()); kafkaTemplate.send(topic, key, payload).whenComplete((metadata, ex) -> { if (ex != null) { - log.error("[ZT-MESSAGE] produce failed: {}", ex.getMessage(), ex); + log.error("[ZT-MESSAGE] produce failed: messageId={}, error={}", msg.getMessageId(), ex.getMessage(), ex); } else if (metadata != null) { - log.info("[ZT-MESSAGE] produced: partition={}, offset={}", metadata.getRecordMetadata().partition(), metadata.getRecordMetadata().offset()); + log.info("[ZT-MESSAGE] produced: messageId={}, partition={}, offset={}", msg.getMessageId(), metadata.getRecordMetadata().partition(), metadata.getRecordMetadata().offset()); } }); }