【发布时间】:2022-01-19 21:15:12
【问题描述】:
由于 flink 问题:https://issues.apache.org/jira/browse/FLINK-24697,我正在尝试手动创建一个消费者组,因此 flink 作业可以成功运行,但是使用以下代码,我无法创建一个,我是否遗漏了什么。
val properties = new Properties()
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker)
properties.put(ConsumerConfig.GROUP_ID_CONFIG, groupId)
properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, resetMode)
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer])
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer])
val consumer = new KafkaConsumer(properties)
consumer.subscribe(util.Arrays.asList(topic))
consumer.close()
日志:
2022-01-19 12:29:55,006 INFO org.apache.flink.avro.registry.confluent.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka version: 2.4.0
2022-01-19 12:29:55,006 INFO org.apache.flink.avro.registry.confluent.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka commitId: 77a89fcf8d7fa018
2022-01-19 12:29:55,006 INFO org.apache.flink.avro.registry.confluent.shaded.org.apache.kafka.common.utils.AppInfoParser [] - Kafka startTimeMs: 1642595395006
2022-01-19 12:29:55,006 INFO org.apache.flink.avro.registry.confluent.shaded.org.apache.kafka.clients.consumer.KafkaConsumer [] - [Consumer clientId=consumer-testCG-7, groupId=testCG] Subscribed to topic(s): testTopic.
有什么方法可以在不消耗任何数据的情况下手动创建CG
【问题讨论】: