微信公众号搜"智元新知"关注
微信扫一扫可直接关注哦!

Spring Cloud Stream 从 Kafka 主题反序列化无效的 JSON

如何解决Spring Cloud Stream 从 Kafka 主题反序列化无效的 JSON

我正在努力将 Spring Cloud Streams 与 Kafka binder 集成。目的是我的应用程序使用主题中的 json 并将其反序列化为 Java 对象。我使用的是函数方法而不是命令式。我的代码使用结构良好的 json 输入。

另一方面,当我发送无效的json时,我希望触发错误记录方法。这在某些测试用例中有效,而在其他测试用例中无效。我的应用程序反序列化 json,即使它无效并触发包含逻辑的方法,而不是错误记录方法

我无法解决框架反序列化一些非结构化json输入的问题。

@Builder
@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
public class KafkaEventRecord {

    @JsonProperty(value = "transport_Metadata",required = true)
    @NonNull
    private TransportMetadata transportMetadata;

    @JsonProperty(value = "payload",required = true)
    @NonNull
    private Payload payload;
}


@Component
public class TokenEventConsumer {

    @Bean
    Consumer<KafkaEventRecord> consumer() {
        return event -> {
            log.info("Kafka Event data consumed from Kafka {}",event);
        };
    }
}

@Configuration
@Slf4j
public class CloudStreamErrorHandler {

    @ServiceActivator(inputChannel = "errorChannel")
    public void handleError(ErrorMessage errorMessage) {
            log.error("Error Message is {}",errorMessage);
    }
}

@EmbeddedKafka(topics = {"batch-in"},partitions = 3)
@TestPropertySource(properties = {"spring.kafka.producer.bootstrap-servers=${spring.embedded.kafka.brokers}","spring.kafka.consumer.bootstrap-servers=${spring.embedded.kafka.brokers}","spring.cloud.stream.kafka.binder.brokers=${spring.embedded.kafka.brokers}"})
@SpringBoottest(webEnvironment = SpringBoottest.WebEnvironment.RANDOM_PORT)
@ActiveProfiles("test")
@Slf4j
public class KafkaTokenConsumerTest {

    private static String TOPIC = "batch-in";

    @Autowired
    private EmbeddedKafkabroker embeddedKafkabroker;

    @Autowired
    private KafkaTemplate<String,String> kafkaTemplate;

    @Autowired
    private KafkaListenerEndpointRegistry endpointRegistry;

    @Autowired
    private ObjectMapper objectMapper;

    @SpyBean
    KafkaEventHandlerFactory kafkaEventHandlerFactory;

    @SpyBean
    CloudStreamErrorHandler cloudStreamErrorHandler;


    @BeforeEach
    void setUp() {
        for (MessageListenerContainer messageListenerContainer : endpointRegistry.getListenerContainers()) {
            ContainerTestUtils.waitForAssignment(messageListenerContainer,embeddedKafkabroker.getPartitionsPerTopic());
        }
    }

    // THIS METHOD PASSES
    @Test
    public void rejectCorruptedMessage() throws ExecutionException,InterruptedException {

        kafkaTemplate.send(TOPIC,"{{{{").get(); // synchronous call

        CountDownLatch latch = new CountDownLatch(1);
        latch.await(5L,TimeUnit.SECONDS);

        // The frame works tries two times,no idea why
        verify(cloudStreamErrorHandler,times(2)).handleError(isA(ErrorMessage.class));

    }

    // THIS METHOD FAILS
    @Test
    public void rejectCorruptedMessage2() throws ExecutionException,"{}}}").get(); // synchronous call

        CountDownLatch latch = new CountDownLatch(1);
        latch.await(5L,times(2)).handleError(isA(ErrorMessage.class));

    }
}
spring.cloud.stream.kafka.bindings.consumer-in-0.consumer.configuration.spring.deserializer.key.delegate.class=org.apache.kafka.common.serialization.StringDeserializer
spring.cloud.stream.kafka.bindings.consumer-in-0.consumer.configuration.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
spring.cloud.stream.kafka.bindings.consumer-in-0.consumer.configuration.key.deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
spring.cloud.stream.kafka.bindings.consumer-in-0.consumer.configuration.value.deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer

// Producer only for testing purpose
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
rejectCorruptedMessage 测试方法中的

json,触发 handleError(ErrorMessage errorMessage) 方法,这是预期的,因为它是无效的 json。另一方面, rejectCorruptedMessage2 测试方法中的 .json 触发 TokenEventConsumer 类中的 Consumer<KafkaEventRecord> consumer() 方法,这不是预期的行为,但是,我得到了具有空值的 KafkaEventRecord 对象。

解决方法

Jackson 不认为这是无效的 JSON,它只是忽略尾随的 }} 并将 {} 解码为空对象。

public class So67804599Application {

    public static void main(String[] args) throws Exception {
        ObjectMapper mapper = new ObjectMapper();
        JavaType type = mapper.constructType(Foo.class);
        Object foo = mapper.readerFor(Foo.class).readValue("{\"bar\":\"baz\"}");
        System.out.println(foo);
        foo = mapper.readerFor(Foo.class).readValue("{}}}");
        System.out.println(foo);
    }

    public static class Foo {

        String bar;

        public String getBar() {
            return this.bar;
        }

        public void setBar(String bar) {
            this.bar = bar;
        }

        @Override
        public String toString() {
            return "Foo [bar=" + this.bar + "]";
        }

    }

}
Foo [bar=baz]
Foo [bar=null]

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。