【问题标题】:how to handle one to many relationship using kafka streams join operations如何使用kafka流连接操作处理一对多关系
【发布时间】:2020-11-04 08:37:37
【问题描述】:

您能帮我如何使用 Kafka 流实现这一点吗?

场景:对订单数据的所有发票进行分组。在实时流媒体中,接收发票可能会有延迟。所以我们要等待 20 分钟,在加入之前对所有发票进行分组。

示例:订单“x”有 3 张发票,预计将在 20 分钟内收到。

预期输出:订单和 3 张发票应作为输出主题中的单个数据提供。

我们有以下拓扑来实现这一点。

  1. 我们分别有订单流和发票流

  2. 我们根据订单键对发票进行分组。我们设置了 20 分钟的翻转窗口

  3. 将订单数据与生成的发票组连接

  4. 将输出写入新主题

问题:第 3 步没有等待第 2 步完成。收到订单后立即加入操作。所以我们没有得到预期的输出。

我们尝试使用连接窗口来达到同样的效果。但由于连接窗口是滑动窗口,我们在输出主题中得到重复数据。

对于上面的例子,如果我们使用连接窗口而不是翻转窗口,我们将得到 3 个输出数据,订单分别有 1 个发票、2 个发票和 3 个发票。

请帮助我解决此问题或建议任何替代方法

代码sn-p:

 KTable<Windowed<String>, List<InvoiceList>> invoiceList= invoiceStream
                .groupByKey()
        .windowedBy(TimeWindows.of(Duration.ofSeconds(1200)))
                .aggregate(() -> new ArrayList<InvoiceList>(),
                        (key, newValue, agg) -> {
                            new KeyValue<>(key, agg.add(newValue));
                            return agg;
                        },
                        Materialized.as("invoice-list").with(Serdes.String(), new ArrayListSerde<InvoiceList>(AppSerdes.InvoiceList())))
 
KStream<String, Order> orderOutput=
 
                orderStream.join(invoiceList, Joiner);
 
       
        orderOutput.to(AppConfig.OutputTopic.OUTPUT_ORDER,Produced.with(Serdes.String(), AppSerdes.Order()));

【问题讨论】:

    标签: apache-kafka apache-kafka-streams kafka-join


    【解决方案1】:

    我假设,订单先到,然后是发票,而不是其他方式。如果我的假设是正确的,那么您的逻辑将不起作用。因为当订单进入您的 KStream 时,可能没有发票,因此连接不会获取任何发票。请记住,KStream-KTable 连接是 non-windowed joins,可以像查找 KTable(更改日志流)一样使用。

    【讨论】:

      【解决方案2】:

      在我们的例子中,这种加入力是有效的。所以我们将它作为两个单独的流接收,并在消费者上添加了自定义逻辑来处理我们的用例。

      谢谢!

      【讨论】:

      • 我建议不要使用普通的KafkaConsumer,而是使用KafkaStreams,但不要使用DSL,您可以使用处理器API:优点是KafkaStreams可以帮助您要获得容错状态,这将很难自己构建。
      猜你喜欢
      • 2016-04-24
      • 2019-04-21
      • 2010-11-20
      • 1970-01-01
      • 1970-01-01
      • 2022-01-27
      • 2022-09-29
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多