SpringBoot中如何实现可变嵌套对象蛇形转驼峰并容错至DLQ?
Hey there! Let's tackle your two key challenges: converting the dynamic details field from snake_case to camelCase, and ensuring failed messages get routed to a DLQ. Here's a straightforward approach tailored for your Spring Boot/Spring Cloud Stream setup:
1. Fixing Snake_Case to CamelCase Conversion for details
The @JsonNaming annotation works great for top-level fields, but it doesn't apply to untyped nested objects like your Object details. Instead, we can use @JsonAnySetter to dynamically capture all incoming fields in details and convert their keys to camelCase on the fly.
Update your EventContent class like this:
import com.fasterxml.jackson.annotation.JsonAnySetter; import com.fasterxml.jackson.databind.PropertyNamingStrategies; import lombok.AllArgsConstructor; import lombok.Getter; import lombok.ToString; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; @ToString @AllArgsConstructor @Getter @JsonNaming(PropertyNamingStrategies.SnakeCaseStrategy.class) // Handles top-level snake_case fields public class EventContent { private final String dataString1; private final String dataString2; // ... other top-level fields private Map<String, Object> details = new HashMap<>(); @JsonAnySetter public void setDetail(String snakeCaseKey, Object value) { // Convert snake_case key to camelCase String camelCaseKey = PropertyNamingStrategies.LOWER_CAMEL_CASE.translate(snakeCaseKey); // Recursively handle nested objects/lists if (value instanceof Map) { Map<String, Object> nestedCamelCaseMap = new HashMap<>(); ((Map<String, Object>) value).forEach((nestedKey, nestedValue) -> nestedCamelCaseMap.put(PropertyNamingStrategies.LOWER_CAMEL_CASE.translate(nestedKey), nestedValue) ); details.put(camelCaseKey, nestedCamelCaseMap); } else if (value instanceof List) { List<Object> nestedCamelCaseList = new ArrayList<>(); ((List<Object>) value).forEach(item -> { if (item instanceof Map) { Map<String, Object> itemMap = new HashMap<>(); ((Map<String, Object>) item).forEach((itemKey, itemValue) -> itemMap.put(PropertyNamingStrategies.LOWER_CAMEL_CASE.translate(itemKey), itemValue) ); nestedCamelCaseList.add(itemMap); } else { nestedCamelCaseList.add(item); } }); details.put(camelCaseKey, nestedCamelCaseList); } else { details.put(camelCaseKey, value); } } }
This setup:
- Uses
@JsonAnySetterto catch all dynamic fields in thedetailssection - Converts every snake_case key to camelCase (including nested keys)
- Handles both nested objects and lists seamlessly
- Keeps your top-level fields covered with the existing
@JsonNamingannotation
2. Setting Up DLQ for Fault Tolerance
Spring Cloud Stream's Kafka binder has built-in support for DLQs, so you don't need to reinvent the wheel. Here's how to configure it:
Add these properties to your application.yml (adjust names to match your setup):
spring: cloud: stream: kafka: binder: brokers: your-kafka-broker-address bindings: # Replace with your input binding name (matches MessagingChannels.BLA_CHANNEL) blaChannel-in-0: consumer: # Enable automatic DLQ routing enable-dlq: true # Name of your DLQ topic dlq-name: your-dlq-topic-name # Optional: Configure retry before sending to DLQ retry-template-name: eventProcessingRetryTemplate dlq-properties: key.serializer: org.apache.kafka.common.serialization.StringSerializer value.serializer: org.springframework.kafka.support.serializer.JsonSerializer # Define your channel bindings if not already done bindings: blaChannel-in-0: destination: your-input-topic-name group: your-consumer-group-name
Then add a retry template bean to handle temporary failures (like MongoDB downtime):
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.retry.backoff.FixedBackOffPolicy; import org.springframework.retry.policy.SimpleRetryPolicy; import org.springframework.retry.support.RetryTemplate; @Configuration public class RetryConfig { @Bean(name = "eventProcessingRetryTemplate") public RetryTemplate eventProcessingRetryTemplate() { RetryTemplate retryTemplate = new RetryTemplate(); // Retry up to 3 times SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(3); // Wait 1 second between retries FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy(); backOffPolicy.setBackOffPeriod(1000); retryTemplate.setRetryPolicy(retryPolicy); retryTemplate.setBackOffPolicy(backOffPolicy); return retryTemplate; } }
What This Does:
- Automatic DLQ Routing: If a message fails deserialization (e.g., invalid JSON) or processing (after retries), it gets sent directly to your configured DLQ topic.
- Retry Logic: Temporary issues (like a brief MongoDB outage) will trigger retries before the message is marked as failed.
- Minimal Code Changes: No need to wrap your listener in try/catch blocks for basic failure handling—Spring handles it for you.
Bonus: Manual DLQ Sending for Custom Errors
If you need to explicitly send a message to DLQ during processing (e.g., business logic validation failure), inject StreamBridge and use it to send messages:
import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; @Component public class EventListener { private final StreamBridge streamBridge; public EventListener(StreamBridge streamBridge) { this.streamBridge = streamBridge; } @StreamListener(MessagingChannels.BLA_CHANNEL) public void onIncomingEvent(EventContent eventContent, @Header(name = "uuid") UUID messageId) { try { doStuffWith(eventContent); } catch (CustomValidationException e) { // Send to DLQ manually streamBridge.send( MessagingChannels.DLQ_CHANNEL, MessageBuilder.withPayload(eventContent) .setHeader("uuid", messageId) .setHeader("errorReason", e.getMessage()) .build() ); } } private void doStuffWith(EventContent eventContent) { // Your business logic + MongoDB persistence } }
内容的提问来源于stack exchange,提问作者user117788

