【问题标题】:Apache Flink Kafka Consumer IssueApache Flink Kafka 消费者问题
【发布时间】:2020-08-08 06:19:34
【问题描述】:

我在 Kafka 中有数据,我想读取 Kafka 是否发送数据的数据,然后过滤它们并返回 JSON。

        // create execution environment
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        Properties properties = new Properties();

        properties.setProperty("bootstrap.servers", "localhost:9092");
       
        properties.setProperty("group.id", "flink_consumer");


        FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>("test-topic",
                new SimpleStringSchema(), properties);
        consumer.setStartFromLatest();
        //config.setWriteTimestampToKafka(true);

        DataStream<String> stream = env.addSource(consumer);

        stream.map(new MapFunction<String, String>() {
            private static final long serialVersionUID = 1L;
            @Override
            public String map(String value) throws Exception {
                
                return "Stream Value: " + value;
            }
        }).print();
        env.execute();

案例 1:当 Kafka 生产者将数据发送到 Kafka 时,我可以在控制台中看到值打印。 - 这很好。 案例 2:Kafka 生产者停止发送数据,Kafka 仍然在主题中具有价值,但相同的代码没有返回任何数据。 -- 这可能吗?

知道哪里出错了吗?

{"firsname":"test", "lastname":"topic", "value":"3.45", "location":"UK"}

我想要过滤 firstname 并返回 JSON。

我看到在数据流处理过程中有过滤器选项。

【问题讨论】:

  • 如果你想从第一条消息开始,你应该设置consumer.setStartFromEarliest();。它将从第一个未确认的消息开始读取。
  • Zahid - 非常感谢,它确实有效。
  • 我很高兴它有帮助。
  • 当然,我看不到点赞按钮。
  • zahid 应该把他的评论变成答案,然后你就可以投票了。

标签: apache-kafka apache-flink flink-streaming


【解决方案1】:

如果你想从第一条消息开始,你应该设置consumer.setStartFromEarliest();。它将从第一个未确认的消息开始读取。

【讨论】:

    猜你喜欢
    • 2017-08-18
    • 1970-01-01
    • 2018-12-31
    • 2018-10-15
    • 1970-01-01
    • 1970-01-01
    • 2022-01-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多