【发布时间】:2020-08-04 07:36:45
【问题描述】:
我们正在创建一个 POC 来读取数据库 CDC 并将其推送到外部系统。
- 每个源表 CDC 以 Avro 格式发送到各自的主题(带有 Kafka Schema Registry 和 Kafka 服务器)
- 我们正在编写 java 代码来使用 avro 模式中的消息,使用 AvroSerde 对其进行反序列化并加入它们,然后发送到不同的主题,以便外部系统可以使用它。
我们有一个限制,尽管我们无法向源表主题生成消息以发送/接收新内容/更改。因此,编写连接代码的唯一方法是在我们运行应用程序时每次从每个源主题从头开始读取消息。(直到我们确信代码正在工作并且可以再次开始接收实时数据)
在 KafkaConsumer 对象中,我们可以选择使用 seekToBeginning 方法强制从 java 代码开始读取,这很有效。但是,当我们尝试使用 KStream 对象流式传输主题并强制从头开始读取它时,没有选项。这里有什么替代方案?
我们尝试使用带有 --to-earliest 的 kafka-consumer-groups reset-topic 重置偏移量,但这只会将偏移量设置为最近的 .当我们尝试使用 --to-offset 参数使用“0”手动重置偏移量时,我们得到低于警告但未设置为“0”。我的理解是,设置为 0 应该从头开始阅读消息。如果我错了,请纠正我。
"WARN 新偏移量 (0) 低于主题分区的最早偏移量"
下面的示例代码
Properties properties = new Properties();
properties.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVER);
properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID);
properties.put("schema.registry.url", SCHEMA_REGISTRY_URL);
properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
properties.put(StreamsConfig.APPLICATION_ID_CONFIG, APPLICATION_ID);
StreamsBuilder builder = new StreamsBuilder();
//nothing returned here, when some offset has already been set
KStream myStream = builder.stream("my-topic-in-avro-schema",ConsumedWith(myKeySerde,myValueSerde));
KafkaStreams streams = new KafkaStreams(builder.build(),properties);
streams.start();
【问题讨论】:
标签: java apache-kafka