【问题标题】:Kafka connect SourceTask commit() and commitRecord() methodsKafka 连接 SourceTask 的 commit() 和 commitRecord() 方法
【发布时间】:2021-05-14 03:15:39
【问题描述】:

我是 Kafka 连接的新手,我正在尝试为我的自定义 JDBC 源连接器(从 oracle DB 读取)构建确认机制。因此,每当数据被添加到 Kafka 主题时,我都想更新源数据库表中的状态/偏移量。 Kafka connect 的融合文档提到了 2 种方法:commit 和 commitRecord,但指出“API 是为具有消息确认机制的源系统提供的”(参考:https://docs.confluent.io/platform/current/connect/devguide.html,请参阅部分:“任务示例 - 源任务")

  1. oracle DB 是否支持确认机制?
  2. 如果可以,我们可以使用 commit() 或 commitRecord() 来更新源 DB 中的状态/偏移量吗?
  3. 如何实现这些方法?
  4. 我们可以为此使用默认的 JDBC 源连接器吗? (https://docs.confluent.io/3.2.0/connect/connect-jdbc/docs/source_connector.html)

【问题讨论】:

  • 您可以完全控制可以读取/写入数据库的内容,所以是的,它支持“确认”。 Confluent 存储最近读取的主 ID 或时间戳

标签: jdbc apache-kafka apache-kafka-connect


【解决方案1】:

我想知道为什么要标记源 Oracle 表中已读取的记录?如果某些内容被写入 Kafka 主题,则意味着它是从源代码中读取的。在这种情况下,您可以将 Confluent 的 JdbcSourceConnectorOracleDatabaseDialect 一起使用。

您当然可以创建 Sink 连接器,该连接器将从主题中读取并更新源表中的记录,但它是为艺术而艺术。

【讨论】:

  • 我想标记/确认源表中的记录,以确保一致性和准确性。实际上,该解决方案甚至不能丢失一行数据。如果Kafka connect从源表读取,由于任何故障而无法写入kafka topic,则需要源系统知道。
  • @JavaDev 如果你不能错过任何事件,那就使用 Debezium 或 GoldenGate
  • @OneCricketeer 我将不得不探索 Debezium 或 GoldenGate。但是根据您的回复,oracle 应该支持确认,我想可以覆盖 commit 和 commitRecord。会尝试的。谢谢!
  • 一个原因可能是实现了许多客户端可以一次写入特定数据库的发件箱模式。然后您不能将递增的 id 用作偏移量(因为事务可能无序提交,因此您可能会丢失一些消息)。我对 Kafka Connect 不太熟悉,但 commit()/commitRecord() 操作似乎是一个希望将发件箱记录标记为已处理的地方。
猜你喜欢
  • 2017-09-03
  • 1970-01-01
  • 2019-11-23
  • 2019-03-10
  • 2021-02-10
  • 2021-02-18
  • 2015-09-03
  • 2019-03-06
  • 2019-12-20
相关资源
最近更新 更多