【发布时间】:2021-07-14 02:14:52
【问题描述】:
我想通过 KTable 加入 KStream。两者都有不同的密钥,但使用自定义分区器共同分区。但是,连接不会产生结果。
KStream 具有以下结构
- 键:房屋 - 组
- 值:用户
KTable 具有以下结构
- 键:用户 - 组
- 值:地址
为了确保每个插入两个主题都按插入顺序进行处理,我使用了一个自定义分区器,我使用每个键的 Group 部分对两个主题进行分区。
我希望得到以下结构的流:
- 键:房屋 - 组
- 值:用户 - 地址
为此,我正在执行以下操作:
val streamsBuilder = streamBuilderHolder.streamsBuilder
val houseToUser = streamsBuilder.stream<HouseGroup, User>("houseToUser")
val userToAddress = streamsBuilder.table<UserGroup, Address>("userToAddress")
val result: KStream<HouseGroup, UserWithAddress> = houseToUser
.map { k: HouseGroup, v: User ->
val newKey = UserGroup(v, k.group)
val newVal = UserHouse(v, k.house)
KeyValue(newKey, newVal)
}
.join(userToAddress) { v1: UserHouse, v2: Address ->
UserHouseWithAddress(v1, v2)
}
.map{k: UserGroup, v: UserHouseWithAddress ->
val newKey = HouseGroup(v.house, k.group)
val newVal = UserWithAddress(k.user, v.address)
KeyValue(newKey, newVal)
}
这需要一个匹配的连接,但它不起作用。
我想显而易见的解决方案是加入一个全局表并放弃自定义分区器。但是,我仍然不明白为什么上述方法不起作用。
【问题讨论】:
标签: java apache-kafka apache-kafka-streams