【发布时间】:2020-07-21 16:22:21
【问题描述】:
我已经看过官方文档了,没办法。
我到 es 的数据来自 kafka,有时可能会出现故障。过去,来自 kafka 的消息被解析并直接插入或更新具有特定 ID 的 ES doc。为了避免旧数据覆盖新数据,我必须检查具有特定 ID 的文档是否已经存在,并且该文档的某些属性是否满足条件。然后我执行 UPDATE 操作(或 INSERT)。
我现在做的是“更新前搜索”。
在更新文档之前,我从 ES 中搜索特定 ID(包含在 kafka msg 中)。然后检查这个文档是否满足条件(例如,update_time是否更旧?)。最后我更新了文档。我将refresh 设置为true 以立即更新索引。
我在担心什么?
似乎是事务性的。
-
如果只有一个线程同步执行,是否有可能在我处理下一条消息时,在上一条消息过程中更新的文档在 ES 中没有刷新?
-
如果我有多个线程消耗 kafka 消息,更新前如何检查?我可以使用脚本来解决这个问题吗?
【问题讨论】:
-
您能解释一下如何让 Kafka 发出“乱序”消息吗?如果与给定文档有关的所有消息都属于同一个分区(它们应该属于),那么您不太可能收到乱序的消息。
-
消息来自多个生产者。
-
是的,但是,所有具有相同 ID 的文档都应该放在同一个分区中(即由文档 ID 键入)
-
但在我的情况下,我使用
process rate之类的东西来识别他们的订单,ID 是一个随机字符串。 -
那你怎么定义文档的ID呢?
标签: go elasticsearch apache-kafka transactions