【发布时间】: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