【问题标题】:Is it possible to "upsert" a message in Kafka using Kafka Connect?是否可以使用 Kafka Connect 在 Kafka 中“更新”消息?
【发布时间】:2018-08-01 19:44:15
【问题描述】:

我正在使用 Confluent 3.3.0。我正在使用jdbc-source-connector 将消息从我的 Oracle 表插入到 Kafka。这工作正常。
我想检查一下“upsert”是否可行。

我的意思是,如果我有一个学生表,有 3 列 id(number)、name(varchar2) 和 last_modified(timestamp)。每当我插入新行时,它将被推送到 Kafka(使用时间戳+自动增量字段)。但是当我更新该行时,应该更新 Kafka 中的相应消息。

我的表的id应该变成他们对应的Kafka消息的key。我的主键 (id) 将作为参考保持不变。
每次更新行时,时间戳字段都会更新。

这可能吗?或者删除 Kafka 中的现有记录并插入新记录。

【问题讨论】:

  • 如果您想要更新事件,您应该考虑使用基于 CDC 的解决方案。 debezium.io/blog/2018/07/19/… 否则,您只是在重新轮询整行,不知道发生了什么变化……而且您也无法在 Kafka 中编辑/删除记录。

标签: jdbc apache-kafka upsert apache-kafka-connect confluent-platform


【解决方案1】:

但是当我更新行时,应该更新 Kafka 中的相应消息

这是不可能的,因为 Kafka 在设计上是只能追加且不可变的。

你会得到的最好的方法是通过某个 last_modified 列查询所有行,或者挂钩 CDC 解决方案,如 Oracle GoldenGate 或 alpha Debezium solution,它将捕获数据库上的单个 UPDATE 事件并附加一个全新的Kafka 主题上的记录。

如果您想对 Kafka 中的数据库记录进行重复数据删除(在某个时间窗口内找到最大 last_modified 的消息),您可以使用 Kafka Streams 或 KSQL 执行该类型的后处理过滤。

如果你使用的是压缩的 Kafka 主题,并且插入了你的数据库键作为 Kafka 消息键,那么在压缩后,最新附加的消息将保持不变,之前具有相同键的消息将被丢弃,而不是更新

【讨论】:

  • 好的。自从我引用 last_modified 列以来,我每次都从 Kafka 获取更新的行。但是,如果我使用与 Kafka 中现有条目相同的消息密钥向 Kafka 中插入一条消息,情况会怎样?会被取代吗?
  • 如前所述,消息始终是附加的,而不是替换的。仅当您已压缩主题时,旧密钥才会在一段时间后被删除
  • 谢谢@cricket_007
猜你喜欢
  • 1970-01-01
  • 2018-03-25
  • 2021-08-06
  • 1970-01-01
  • 1970-01-01
  • 2020-03-02
  • 1970-01-01
  • 2022-10-24
  • 2020-10-08
相关资源
最近更新 更多