问题描述
我正在努力将 Spring Cloud Streams 与 Kafka binder 集成。目的是我的应用程序使用主题中的 json 并将其反序列化为 Java 对象。我使用的是函数式方法而不是命令式。我的代码使用结构良好的 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]