spring cloud stream kafka streams DLQ

Viewed 487

I am using Apache Kafka 2.7.0 with Spring Cloud Stream Kafka Streams.

In my Spring Cloud Stream (Kafka Streams) application, I have configured my application.yml to use the sendToDlq mechanism when the messages in the input topic have deserialization errors :

spring:
  cloud:
    stream:
      function:
        definition: processor      
      bindings:         
        processor-in-0:
          destination: input-topic
          consumer:
            dlqName: input-topic-dlq
        processor-out-0:
          destination: output-topic       
      kafka:
        streams:
            binder:
              deserialization-exception-handler: sendToDlq
            configuration:
              metrics.recording.level: DEBUG
            brokers:
              - localhost:9092

I start my application and I do not see this topic existing. The documentation states that the DLQ topic will be created if not present.

If I try to consume from the DLQ topic, I get an error as below :

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic input-topic-dlq --property print.value=true --property print.key=true --from-beginning
[2021-03-19 10:17:09,936] WARN [Consumer clientId=consumer-console-consumer-85295-1, groupId=console-consumer-85295] Error while fetching metadata with correlation id 2 : {input-topic-dlq=LEADER_NOT_AVAILABLE} (org.apache.kafka.clients.NetworkClient)

At this instant, when I query Zookeeper ls /brokers/topics then I see the Topic created.

Now, I try to POST a non-JSON message to the input-topic ( My default Deserializer is JSON).

BUT I cannot see any messages in the input-topic-dlq topic created.

What is strange is that I can see messages in the default "error.input-topic-dlq.appId" topic.

Am I doing something wrong here ?

1 Answers

I managed to figure that out. There seems to be a typo in current documentation for Spring Cloud Stream Kafka Streams Binder.

Destination for binding should be on spring.cloud.streams.bindings level ike you already have but consumer properties specific to the implementation should be on spring.cloud.streams.kafka.streams.bindings level.

So your configuration should look like this:

spring:
  cloud:
    stream:
      function:
        definition: processor      
      bindings:         
        processor-in-0:
          destination: input-topic
        processor-out-0:
          destination: output-topic       
      kafka:
        streams:
            binder:
              deserialization-exception-handler: sendToDlq
            bindings:
              processor-in-0:
                consumer:
                  dlqName: input-topic-dlq
            configuration:
              metrics.recording.level: DEBUG
            brokers:
              - localhost:9092
Related