【问题标题】:Kafka Streams KGroupedTable.count() returning negative value. How's that possible?Kafka Streams KGroupedTable.count() 返回负值。这怎么可能?
【发布时间】:2019-02-09 11:59:01
【问题描述】:

KGroupedTable.count() 返回负值?

idAndJobTransaction
                .filter((k,v) -> v!=null)
                .mapValues(jobTransaction -> {
                    jobTransaction.setCount(0);
                    jobTransaction.setId(0L);
                    jobTransaction.setRunsheet_id(0L);
                    jobTransaction.setTimestamp(0L);
                    if(jobTransaction.getDelete_flag() == 1)
                        return null;
                    else
                        return jobTransaction;
                } )
                .groupBy((id,jobTransaction)->new KeyValue<>(jobTransaction,jobTransaction),Serialized.with(jobTransactionSerde,jobTransactionSerde))
                .count()
                .toStream()
                .mapValues((k,v)-> new JobSummary(k,v))
                .peek((k,v)->{
                    log.info(k.toString());
                    log.info(v.toString());
                }).selectKey((k,v)-> v.getCompany_id())  // So that the count is consumed in order for each company
                .to(JOB_SUMMARY,Produced.with(Serdes.Long(),jobSummarySerde));

count 方法有时会返回负值。大约 1% 的值是负数。这怎么可能?

编辑 1:

我将此聚合的结果推送到 Postgres 表。负值不限于 -1,但它会达到非常高的值。

我正在使用 2 个消费者。这有什么不同吗?

这可能是 Kafka 流的问题吗?还是我应该看看其他可能的原因?

编辑 3: 我能够捕获一些可用的日志,并且确实在 peek 中看到了负值:

至于 JobSummary 类,它确实是一个非常简单的 POJO 类。这是在 KStream 应用程序中调用的构造函数。

  public JobSummary(JobTransaction j, Long count){
    this.setUser_id(j.getUser_id());
    this.setHub_id(j.getHub_id());
    this.setCity_id(j.getCity_id());
    this.setCompany_id(j.getCompany_id());
    this.setJob_master_id(j.getJob_master_id());
    this.setJob_status_id(j.getJob_status_id());
    this.setCount(count);
    this.setDate(j.getDate());
}

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    我猜(这是我能想到的唯一解释),这是一个特殊的极端情况。首先,您必须了解 KTable 聚合在内部是如何工作的。这是在另一个问题上解释的:TopologyTestDriver sending incorrect message on KTable aggregations

    在这种背景下,如果结果表中的当前计数为零,并且上游基表(即idAndJobTransaction)获得幂等更新(即基表中的记录),则可能发生负计数。表从&lt;K,V&gt; 更新到&lt;K,V&gt;。这将导致一个减法和一个加法记录转到结果表中的同一行(请注意,Kafka Streams 不会比较表更新时的新旧值和盲目假设两者不同)。此外,减法和加法记录是独立发送到下游的,下游count()分两步更新其结果。因此,结果表中的计数从0变为-1处理减法记录和从 -1 回到 0 处理添加记录。

    【讨论】:

    • 但最终的结果应该总是大于或等于零。不应该吗?
    • 我在问题本身中添加了更多细节。
    • 是的,它应该始终归零。我不知道 atm,你怎么会得到小于 -1 的计数......你能通过 count().toStream().print() 确认问题出在count() - 也许它在计数后搞砸了。
    • 这可能是一个 serde 问题吗?
    • @MatthiasJ.Sax 我将添加其他日志并回复您。机会较少,因为在计算计数后没有处理,它只是转储到 Postgres 中。我将添加日志并回复您。因为这个周末我不在城里,所以我会在星期一之前回复你。
    猜你喜欢
    • 1970-01-01
    • 2011-05-11
    • 2017-06-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-09
    • 2011-08-11
    • 2011-01-31
    相关资源
    最近更新 更多