You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 @JsonAnySetter to catch all dynamic fields in the details section
  • 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 @JsonNaming annotation

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 08:06:37