【问题标题】:Kafka Streams wait function with depending objectsKafka Streams等待函数与依赖对象
【发布时间】:2018-02-05 11:42:49
【问题描述】:

我创建了一个 Kafka Streams 应用程序,它接收来自不同主题的不同 JSON 对象,我想实现某种等待功能,但我不确定如何最好地实现它。

为了简化问题,我将在下一节中使用简化的实体,希望可以很好地描述问题。 所以在我的一个流中,我收到汽车对象,每辆车都有一个 ID。在第二个流中,我接收人员对象,每个人也有一个汽车 ID,并分配给具有此 ID 的汽车。

我想使用我的 Kafka Streams 应用程序从两个输入流(主题)中读取数据,并使用具有相同汽车 ID 的四个人来丰富汽车对象。只有当所有四个人都包含在汽车对象中时,汽车对象才应转发到下一个下游处理器。

我计划为汽车和人对象创建一个输入流,将 JSON 数据解析为内部对象表示,将两个流合并在一起,并在合并的流上应用“selectKey”函数来提取键出实体。 之后,我会将数据推送到包含状态存储的自定义转换函数中。在这个转换函数中,我会将每个到达的汽车对象及其 id 存储在状态存储中。一旦新的人员对象到达,我会将它们添加到状态存储中的相应汽车对象(请忽略此处迟到汽车的情况)。一旦有四个人在汽车对象中,我就会将该对象转发到下一个流函数并将汽车对象从状态存储中删除。

这是一个合适的方法吗?我不确定可伸缩性,因为我必须确保在运行多个实例时,具有相同 id 的 car 和 person 对象将由同一个应用程序实例处理。我会为此使用 selectKey 函数,这样行吗?

谢谢!

【问题讨论】:

    标签: stream apache-kafka apache-kafka-streams


    【解决方案1】:

    基本设计对我来说看起来不错。

    但是,selectKey() 本身是不够的,因为transform()(与 DSL 运算符相反)不会触发自动重新平衡。因此,您需要通过through() 手动重新平衡。

    stream.selectKey(...)
          .through("user-created-topic")
          .transform(...);
    

    https://docs.confluent.io/current/streams/upgrade-guide.html#auto-repartitioning

    【讨论】:

    • 谢谢你帮了我很多。我会这样尝试。
    猜你喜欢
    • 2018-08-23
    • 1970-01-01
    • 2020-09-08
    • 1970-01-01
    • 2019-10-18
    • 1970-01-01
    • 1970-01-01
    • 2020-08-15
    相关资源
    最近更新 更多