【发布时间】:2020-12-10 04:05:30
【问题描述】:
我正在开发一个使用 KafkaIO 作为输入的 Beam 应用程序
KafkaIO.<Long, GenericRecord>read()
.withBootstrapServers("bootstrapServers")
.withTopic("topicName")
.withConsumerConfigUpdates(confs)
.withKeyDeserializer(LongDeserializer.class)
.withValueDeserializer((Deserializer.class)
.commitOffsetsInFinalize()
.withoutMetadata();
我试图了解commitOffsetsInFinalize() 的工作原理。
如何完成流式传输作业?
管道的最后一步是自定义 DoFn,它将消息写入DynamoDb。有什么方法可以手动调用一些finalize() 方法,以便在每次成功执行DoFn 后提交偏移量?
我也很难理解检查点和最终确定之间的关系是什么?如果管道上没有启用检查点,我是否仍然能够完成并让commitOffsetsInFinalize() 工作?
p.s 管道现在的方式,即使使用commitOffsetsInFinalize() 读取每条消息,无论下游是否有故障正在提交,因此导致数据丢失。
谢谢!
【问题讨论】:
标签: apache-flink apache-beam apache-beam-kafkaio