【发布时间】:2018-07-27 21:20:42
【问题描述】:
我有一个由其他流组成的流
final KTable<Long, CompositeInfo> compositeInfoTable = compositeImcTable
.leftJoin(
compositeFundTable,
(CompositeImc cimc, CompositeFund cf) -> {
CompositeInfo newCandidate = new CompositeInfo();
if (cimc != null) {
newCandidate.imcName = cimc.imcName;
newCandidate.imcID = cimc.imcID;
if (cf != null) {
newCandidate.investments = cf.investments;
}
}
return newCandidate;
})
.leftJoin(
compositeGeographyTable,
(CompositeInfo cinfo, CompositeGeography cg) -> {
if (cg != null) {
cinfo.regions = cg.regions;
}
return cinfo;
})
.leftJoin(
compositeSectorTable,
(CompositeInfo cinfo, CompositeSector cs) -> {
if (cs != null) {
cinfo.sectors = cs.sectors;
}
return cinfo;
})
.leftJoin(
compositeClusterTable,
(CompositeInfo cinfo, CustomCluster cc) -> {
if (cc != null && cc.clusters != null) {
cinfo.clusters = cc.clusters;
}
return cinfo;
})
.leftJoin(
compositeAlphaClusterTable,
(CompositeInfo cinfo, CompositeAlphaCluster cac) -> {
if (cac != null) {
cinfo.alphaClusters = cac.alphaClusters;
};
return cinfo;
},
Materialized.<Long, CompositeInfo, KeyValueStore<Bytes, byte[]>>as(this.storeName)
.withKeySerde(Serdes.Long())
.withValueSerde(compositeInfoSerde));
我的问题与 CompositeInfo 和 CustomCluster 之间的左连接有关。 CustomCluster 如下所示
KTable<Long, CustomCluster> compositeClusterTable = builder
.stream(
SUB_TOPIC_COMPOSITE_CLUSTER,
Consumed.with(Serdes.Long(), compositeClusterSerde))
.filter((k, v) -> v.clusters != null)
.groupByKey(Serialized.with(Serdes.Long(), compositeClusterSerde))
.reduce((aggValue, newValue) -> newValue);
自定义集群中的消息看起来像
CustomCluster [clusterId=null, clusterName=null, compositeId=280, operation=null, clusters=[Cluster [clusterId=6041, clusterName=MyName]]]
所以我将这个对象中的 HashMap 集群分配给 CompositeInfo 对象中加入到compositeId 上的集群。
我所看到的是,一个给定的compositeId 的CustomCluster 消息进来了一个dis 正确处理,但是包含前一个集群的旧消息(我仍在调查这个)被再次处理。 在挖掘问题发生在kafka内部KTableKTableRightJoin
public void process(final K key, final Change<V1> change) {
// we do join iff keys are equal, thus, if key is null we cannot join and just ignore the record
if (key == null) {
return;
}
final R newValue;
R oldValue = null;
final V2 value2 = valueGetter.get(key);
if (value2 == null) {
return;
}
newValue = joiner.apply(change.newValue, value2);
if (sendOldValues) {
oldValue = joiner.apply(change.oldValue, value2);
}
context().forward(key, new Change<>(newValue, oldValue));
}
当joine第一次返回时,newValue正确更新。但是代码然后转到 sendOldValues 块,一旦加入者返回,newValue 就是更新增益,但这次是旧集群值。
所以这是我的问题:
- 为什么当 joiner 被调用时 newValues 会被更新 第二次使用 oldValue
- 有没有办法关闭 sendOldValues
- 我的链式左连接是否与它有关。我知道 以前版本的 kafka 存在链接错误。但现在我在 1.0
更新: 我发现的另一件事。如果我将连接移动到连接链上并删除其他连接,则 sendOldValues 仍然为 False。因此,如果我有类似以下内容:
final KTable<Long, CompositeInfo> compositeInfoTable = compositeImcTable
.leftJoin(
compositeFundTable,
(CompositeImc cimc, CompositeFund cf) -> {
CompositeInfo newCandidate = new CompositeInfo();
if (cimc != null) {
newCandidate.imcName = cimc.imcName;
newCandidate.imcID = cimc.imcID;
if (cf != null) {
newCandidate.investments = cf.investments;
}
}
return newCandidate;
})
.leftJoin(
compositeClusterTable,
(CompositeInfo cinfo, CustomCluster cc) -> {
if (cc != null && cc.clusters != null) {
cinfo.clusters = cc.clusters;
}
return cinfo;
},
Materialized.<Long, CompositeInfo, KeyValueStore<Bytes, byte[]>>as(this.storeName)
.withKeySerde(Serdes.Long())
.withValueSerde(compositeInfoSerde));
这给了我正确的结果。但我认为,如果我在此之后放置更多的链式连接,它们可能会显示相同的错误行为。
我目前还不确定,但我认为我的问题在于 chained leftjoin 和计算 oldValue 的行为。有没有其他人遇到过这个问题?
更新
经过大量挖掘,我意识到 sendOldValues 是 kafka 内部的,而不是我遇到的问题的原因。我的问题是,当 oldValue 的 ValueJoiner 返回时 newValue 会发生变化,我不知道它是否是由于对 Java 对象的引用分配传递所致
这是传入对象的样子
CustomCluster [clusterId=null, clusterName=null, compositeId=280, operation=null, clusters=[Cluster [clusterId=6041, clusterName=Sunil 2]]]
集群是HashSet<Cluster> clusters = new HashSet<Cluster>();
然后它被连接到一个对象
CompositeInfo [compositeName=BUCKET_NM-280, compositeID=280, imcID=19651, regions=null, sectors=null, clusters=[]]
这里的簇是同类型的,但是在 CompositeInfo 类中
当我加入时,我将 CustomCluster 对象的集群分配给 CompositeInfo 对象
(CompositeInfo cinfo, CustomCluster cc) -> {
if (cc != null && cc.clusters != null) {
cinfo.clusters = cc.clusters;
}
return cinfo;
}
【问题讨论】:
-
“但是代码然后转到 sendOldValues 块,一旦加入者返回,newValue 是更新增益,但这次是旧集群值”是什么意思。
if (sendOldValue)语句设置oldValue =但不设置newValue =。 -
“有没有办法关闭 sendOldValues”。并非如此,它很可能会导致不正确的结果——
sendOldValues用于保证正确的计算。 -
KTable#leftJoin()需要sendOldValues才能计算出正确的结果。 -
@MatthiasJ.Sax 请查看最新更新。我意识到我的问题不是 sendOldValues。我有一个包含 HashMap 的对象。这包括新的和旧的价值观。当 oldValue 代码块的连接器返回时,它会以某种方式更改 newValue 中的 HashMap
-
我解决了。这确实是一个参考传递问题。加入时,我需要初始化并返回一个新对象,而不是给旧对象赋值
标签: java apache-kafka left-join apache-kafka-streams