【问题标题】:commitOffsetsInFinalize() and checkmarks in Apache BeamApache Beam 中的 commitOffsetsInFinalize() 和复选标记
【发布时间】: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


    【解决方案1】:

    这里的 finalize 是指检查点的最终确定,换句话说,当数据被持久地提交到 Beam 的运行时状态时(这样将重试工作失败/重新分配,而无需再次从 Kafka 读取此消息)。这并不意味着有问题的数据已经完成了管道的其余部分。

    【讨论】:

    • 谢谢罗伯特。如果从未定义检查点会发生什么?如果应用在 Flink 上运行,默认的检查点机制是什么?另外我的理解是,除非我有ENABLE_AUTO_COMMIT_CONFIG, false,否则无论是否设置commitOffsetsInFinalize(),消息仍然会被提交,对吗?
    • 正确,这仅在您还 AUTO_COMMIT 未设置 kafka 消费者配置时才有用。
    猜你喜欢
    • 2022-12-24
    • 2018-12-07
    • 1970-01-01
    • 2019-04-25
    • 1970-01-01
    • 1970-01-01
    • 2017-08-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多