spring-boot - 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, 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));
}
}
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 ,triggershandleError(ErrorMessage errorMessage)
方法,这是预期的,因为它是无效的json。另一方面,rejectCorruptedMessage2测试方法中的 json 触发Consumer<KafkaEventRecord> consumer()
了 TokenEventConsumer 类中的方法,这不是预期的行为,但是,我得到了具有空值的 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]
推荐阅读
- java - 后台网络 I/O 的服务、IntentService 或线程
- workflow - 失败后重新启动 Informatica 会话,尝试 N 次或通知
- php - 尽管看到该字段,但使用 Webdriver $I->fillField 的代码接收不起作用
- databricks - 我可以在数据块中创建等效的 SQL 临时表吗?
- javascript - 为什么当我尝试读取 object[key].value 时出现未定义?
- python - 查找给定角色中的所有用户 discord.py
- xamarin.forms - 如何从 ParseObject 的子类中的相关用户获取数据?
- ansible - 需要帮助循环通过 Ansible msg 输出
- swift - 如何在 Combine 中最好地创建 @Published 值的发布者聚合?
- python - 如何在 Python 中解压缩复合文件