【问题标题】:Relative amount of events in kafka streams using join使用连接的 kafka 流中的相对事件数量
【发布时间】:2022-07-04 17:47:11
【问题描述】:

我正在编写一个 kafka 流应用程序,在其中我正在为网页生成统计信息。 我有一个关于网页的信息流,其中包括结构中的页面类型(新闻、游戏、博客等)和页面语言(en、fr、ru 等)。

我已将此流过滤为第二个流,其中包含特定页面类型的所有语言。 对于这个例子,我们可以假设过滤后的流包括“新闻”页面的所有事件。

我现在想将每种语言的页面数量除以相同类型的页面总数的值 a 输出到主题。

我使用 .count() 创建了一个 KTable 来计算每种语言的事件。 我还使用 .count() 创建了一个包含所有相同类型事件的 KTable。

为了产生除法,我计划在流之间使用连接,它将取左值并将其除以右值。 不幸的是,这似乎不起作用,因为左值的键是语言,而右值的键是页面类型。

我的代码如下:

ValueJoiner<Long, Long, Float> valueJoiner = (leftVal, rightVal) -> {
            if ((rightVal != null) && (leftVal != null))
            {        
                return leftVal.floatValue()/rightVal;
            }
            return 0f;
        };

// the per language table for news pages
KTable<String, Long> langTable = newsStream.selectKey((ignored, value) -> value.getLang()).groupByKey().count();
// the table which counts all events of news pages
KTable<String, Long> allTable = newsStream.groupBy((ignored, value) -> value.getType()).count();

// this is the join that doesn't produce values (as there are no common keys?)
KTable<String, Float> joinedLangs = langTable.join(allTable, valueJoiner);

使此代码工作并产生相对数量值的最佳方法是什么?

【问题讨论】:

    标签: apache-kafka kafka-join


    【解决方案1】:

    如果我们在谈论Join,那么两边(左右)的输入数据必须是共同分区的。参考https://developer.confluent.io/tutorials/foreign-key-joins/kstreams.html

    利用用户提供的 ValueJoiner,如下有效地创建联接输出记录:

    KeyValue<K, LV> leftRecord = ...;
    KeyValue<K, RV> rightRecord = ...;
    ValueJoiner<LV, RV, JV> joiner = ...;
    
    KeyValue<K, JV> joinOutputRecord = KeyValue.pair(
        leftRecord.key, /* by definition, leftRecord.key == rightRecord.key */
        joiner.apply(leftRecord.value, rightRecord.value)
      );
    

    加入共分区要求。

    • 加入时输入数据必须是同分区的。
    • 这可确保在处理期间将来自连接两侧的具有相同键的输入记录传递到相同的流任务。
    • 在加入时确保数据共同分区是用户的责任。 考虑使用全局表 (GlobalKTable) 进行连接,因为它们不需要数据共同分区。

    数据共分区的要求是:

    • join 的输入主题(左侧和右侧)必须具有相同的分区数。
    • 写入输入主题的所有应用程序必须具有相同的分区策略,以便将具有相同键的记录传递到相同的分区号。
    • 换句话说,输入数据的键空间必须以相同的方式跨分区分布。 使用 Kafka 的 Java Producer API 的应用程序必须使用相同的分区器(参见生产者设置“partitioner.class”又名 ProducerConfig.PARTITIONER_CLASS_CONFIG),并且使用 - Kafka 的 Streams API 的应用程序必须使用相同的 StreamPartitioner 进行操作,例如 KStream#to ()。
    • 如果您碰巧在所有应用程序中使用默认的分区程序相关设置,则无需担心分区策略。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-08-09
      • 1970-01-01
      • 2015-07-17
      • 2019-03-10
      • 2018-10-31
      • 1970-01-01
      • 1970-01-01
      • 2020-10-05
      相关资源
      最近更新 更多