【问题标题】:kafka streams - all group members handle a global eventkafka 流 - 所有组成员处理一个全局事件
【发布时间】:2020-10-23 12:39:21
【问题描述】:

我的问题是:我应该如何构建我的 kafka 应用程序,以应对使用另一个具有多个分区的主题的单个分区中发生的“全局”事件?

如果我遇到以下情况: 我有一个gigs 主题,每次在一个国家宣布新演出时都会收到一条新消息,key 是艺术家的名字,value 是带有日期、位置等的 avro 对象。主题有 1 分区。

示例键:Abba 示例值:{"date": "2020-08-21T09:00:00Z", "venue":"The O2 London"}

我还有另一个名为 subscribes 的压缩主题,其中包含订阅艺术家的用户,他们希望在他们感兴趣的艺术家宣布演出时收到通知。这里的关键是 avro 密钥以 userId 和艺术家姓名作为值,该值是一个带有电子邮件和其他首选项的 avro 对象。该主题有 10 个分区。

示例键:{"user": 123, "artist": "Abba"} 示例值:{"email": "123@email.com", ...}

我有很多用户,但没有那么多演出,我想运行 10 个应用程序实例,因为这是 subscribes 主题的分区数。

应用程序使用这两个主题,当每个用户在subscribes 主题中有一条消息的gigs 生成一个演出时,key.artist 等于来自gigs 的密钥生成一条消息到第三个主题。我可以通过扫描由subscribes 构建的本地商店来实现这一点。

有人会如何去使用 kafka 流(或者只是消费者)来开发这样的应用程序?该文档提到了全局表作为向所有成员广播事件的一种方式,但我遇到了不同的问题:

  1. 如果我使用来自subscribes 的KTable 和来自gigs 的GlobalTable 编写此应用程序,那么我无法使用dsl 加入这两者。用户可以在subscribes 中停留未知时间,因此我无法使用 KStream 加入 GlobalTable。
  2. 当使用新消息时,使用处理器 api 我无法从全局上下文(来自 gigs 的全局表)转发到本地上下文(我可以从 subscribes 访问本地可查询存储)。
  3. 即使我可以从全局上下文转发到 2 分之一的本地上下文。我遇到的问题是,每当应用程序重新启动时,它都会从头开始读取主题,因为全局表不保留偏移量,所以我会最终将相同的消息一次又一次地发送到输出主题。

我想到的解决方案: 使用 gigs 主题并生成到具有 10 个分区的 gigs_broadcast 并将相同的消息写入每个分区。这样,每个成员都可以从每个存储中分配相同的分区,因此每当一条消息到达gigs 时,它现在都在本地上下文中,并且可以将其转发到另一个处理器,该处理器可以从@987654340 访问可查询的存储@ 进行扫描。

有没有更好的解决方案?也许通过使用消费者 + 生产者而不是 kafka 流?理想情况下,我知道我对这两个主题都有相同的键,但我不确定这是否可行:我必须将subscribes 流重新设置为艺术家,执行 groupBy,聚合每个艺术家的 userIds 列表然后加入将 gigs 主题转换为对它们进行平面映射以生成 gigs_customers ,其中键与 subscribes 相同,值相同。

【问题讨论】:

  • 我认为您所描述的可能是最好的方法:从订阅中消费并立即将每个订阅者放入可查询的状态存储中。从演出中消费并在商店中查找订阅者并转发到第三个主题。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

对我来说,这听起来像是外键连接。您可以将这两个主题读取为 KTables,并将艺术家用作连接属性。您的用户/订阅表将是您的左侧表,并且您使用键提取器作为连接属性提取到艺术家并连接到 gigs 表。

【讨论】:

    猜你喜欢
    • 2014-08-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-03-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多