【问题标题】:Apache Beam, KafkaIO at least once semanticsApache Beam、KafkaIO 至少一次语义
【发布时间】:2021-01-29 08:26:41
【问题描述】:

我们正在实施一个从 Kafka 读取并写入 BigQuery 的试点。

简单的管道:

  • KafkaIO.read
  • BigQueryIO.write

我们关闭了自动提交。 我们正在使用commitOffsetsInFinalize()

如果 BigQueryIO 端一切正常,此设置能否保证消息在 BigQuery 中至少出现一次并且不会丢失?

commitOffsetsInFinalize() 的文档中,我遇到了以下情况:

它有助于在从头开始重新启动管道时最大限度地减少记录的间隙或重复处理

我很好奇这里指的是什么“差距”?

如果考虑边缘情况,是否有可能会跳过消息而不将消息传递到 BQ?

【问题讨论】:

    标签: apache-kafka kafka-consumer-api apache-beam apache-beam-kafkaio


    【解决方案1】:

    提交 Apache Kafka 的偏移量意味着如果您要重新启动管道,它将在您重新启动之前在流中的位置开始。 Dataflow 确实保证在写入 BigQuery 时不会丢弃数据。但是,使用分布式系统时,总是有可能出现问题(例如,GCP 中断)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多