Kafka consuming the latest message again when I rerun the Flink consumer

Viewed 2732

I have created a Kafka consumer in Apache Flink API written in Scala. Whenever I pass some messages from a topic, it duly is receiving them. However, when I restart the consumer, instead of receiving the new or unconsumed messages, it consumes the latest message that was sent to that topic.

Here's what I am doing:

  1. Running the producer:

    $ bin/kafka-console-producer.sh --broker-list localhost:9092 --topic corr2
    
  2. Running the consumer:

    val properties = new Properties()
    properties.setProperty("bootstrap.servers", "localhost:9092")
    properties.setProperty("zookeeper.connect", "localhost:2181")
    properties.setProperty("group.id", "test")
    
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    val st = env
        .addSource(new FlinkKafkaConsumer09[String]("corr2", new SimpleStringSchema(), properties))
    env.enableCheckpointing(5000)
    st.print()
    env.execute()
    
  3. Passing some messages

  4. Stopping the consumer
  5. Running the consumer again prints the last message I sent. I want it to print only new messages.
1 Answers
Related