Avro SerializationException is making Kafka Streams to shutdown causing it to not process next incoming records

Viewed 295

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.

0 Answers
Related