I'm using spring kafka streams in an api which is supposed to consume events from Topic A, transform them and publish them to Topic B.
We are using Avro pojo classes shown below which are converted to Avro message by this property
valueSerde: io.confluent.kafka.streams.serdes.avro.ReflectionAvroSerde.
So upon conversion the Avro message should match with the avro schema registered in the schema registry.
import lombok.Getter;
import lombok.Setter;
import org.apache.avro.reflect.AvroDoc;
import org.apache.avro.reflect.Nullable;
@Getter
@Setter
public class TestEventAttributes {
@AvroDoc("signOnType")
private String signOnType;
@AvroDoc("signOnTypeDescription")
private String signOnTypeDescription;
@Nullable
@AvroDoc("sessionId")
private String sessionId;
}
As you can see some attributes are nullable attributes and others are mandatory. Pretty standard Kafka use case and it is working well too. But there is one error scenario I'm unable to figure out how to handle it. When the attribute which is not nullable is set to null I get the below exception:
org.apache.kafka.streams.errors.StreamsException: Exception caught in process. taskId=0_5, processor=KSTREAM-SOURCE-0000000000, topic=test-events-in, partition=5, offset=141, stacktrace=org.apache.kafka.common.errors.SerializationException: Error serializing Avro message
Caused by: java.lang.NullPointerException: in TestEventAttributes in string null of string in field signOnTypeDescription of com.test.TestEventAttributes
Caused by: org.apache.kafka.common.errors.SerializationException: Error serializing Avro message
When this happens the stream shuts down and it does not process the next incoming record.
The CustomProductionExceptionHandler implementation does not help as it looks like Avro serialization exception is not covered in them.
Sure I can do the null check on each mandatory attribute but when we are dealing with a number of large model classes and thousands of incoming messages which I don't have control over there is a risk of one of the incoming message can be a 'poison pill' message.
I would like to know if there is an elegant way to handle this error scenario and let the stream continue to process the next incoming records.
Here are some relevant portions of the config:
spring:
cloud:
stream:
function.definition: process
schemaRegistryClient:
endpoint: http://test.com:8081
bindings:
process-in-0:
destination: test-events-in
process-out-0:
destination: test-events-out
content-type: application/*+avro
kafka:
streams:
binder:
autoCreateTopics: false
brokers: test:9092, test1:9092
functions:
process:
applicationId: test-events-out-cg1
configuration:
default:
key:
serde: org.apache.kafka.common.serialization.Serdes$StringSerde
value:
serde: io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
deserialization:
exception:
handler : org.apache.kafka.streams.errors.LogAndContinueExceptionHandler
production:
exception:
handler: com.test.events.stream.exception.CustomProductionExceptionHandler
value:
subject:
name:
strategy: io.confluent.kafka.serializers.subject.TopicRecordNameStrategy
bindings.process-out-0.producer:
keySerde: org.apache.kafka.common.serialization.Serdes$StringSerde
valueSerde: io.confluent.kafka.streams.serdes.avro.ReflectionAvroSerde
schema.registry.url: http://test.com:8081
useNativeEncoding: true
Others have reported on the same problem here and here but no reported solution yet.
Also there is this jira issue which is still open.