基于Kafka Streams按Key分组并创建键为字符串值为对象列表的全局KTable
数据表
| SCHEME | RULEORDER | REGEX |
|---|---|---|
| VS | 1 | 0*[2368A] |
| MC | 1 | 0*[23] |
| MC | 2 | 0* |
| ZS | 1 | 9* |
| MC | 3 | 22* |
数据发送代码说明
数据表的每一行数据都会以KStream形式发送,发送代码如下:
ruleStream .peek((k, v) -> logReceivedRecord("RULE", k, Optional.ofNullable(v).map(RuleCdc::getUPDATETS).orElse(null))) .filter((k, v) -> Objects.nonNull(v), Named.as("rule-nonnull-filter")) .map((k, v) -> new KeyValue<>(v.getSCHEME()+ "-"+ v.getRULEORDER(), mapper.mapconfig(v, initializer.initConfig(k)))) .peek((k, v) -> logSentRecord("RuleConfig", k, getTopicNameWithRegion(TOPIC_NAME)));
这里用v.getSCHEME() + "-" + v.getRULEORDER()作为Key,目的是后续创建KTable时保留每个SCHEME对应的所有RULEORDER。如果仅用v.getSCHEME()做Key,同一SCHEME下的最新RULEORDER会覆盖旧值。
需求与问题
我需要将KStream中逐条接收的数据按SCHEME作为Key聚合,让每个SCHEME对应一个包含所有关联RULEORDER的列表。但自己尝试的代码未得到预期结果。
尝试的代码
StreamsBuilder builder = new StreamsBuilder(); KStream<String , RuleConfig> ruleConfigKStream = builder.stream(TOPIC_NAME,Consumed.with(stringSerde, ruleConfigSerde)); KGroupedStream<String , RuleConfig> groupedKStream = ruleConfigKStream.groupBy((key, value) -> value.getScheme(), Grouped.with(Serdes.String(),ruleConfigSerde)); KTable<String,List<RuleConfig>> ruleStore = groupedKStream.aggregate(()-> new ArrayList<>(), (key,value,list) -> { list.add(value) ; Materialized.<String, List<RuleConfig>, KeyValueStore<Bytes,byte[]>> as (RULE_STORE) .withKeySerde(stringSerde).withValueSerde(listSerde); return list; }); builder.build(); ReadOnlyKeyValueStore<String, List<RuleConfig>> ruleKVStore = kafkaStreams.store(StoreQueryParameters.fromNameAndType(RULE_STORE, QueryableStoreTypes.keyValueStore())); ruleKVStore.get("MC"); --> **未得到预期结果**
问题分析与修正方案
你的代码存在3个关键问题:
Materialized配置位置错误:它应该作为aggregate方法的第三个参数传入,而非放在聚合逻辑的lambda内部,当前写法会导致配置无法正确生效。- 可变列表的状态更新问题:直接修改传入的列表并返回,Kafka Streams的状态存储无法正确处理这种可变对象的更新,建议每次聚合创建新列表。
- 可选的排序逻辑:如果需要列表按RULEORDER有序,需手动添加排序步骤。
修正后的代码:
StreamsBuilder builder = new StreamsBuilder(); KStream<String, RuleConfig> ruleConfigKStream = builder.stream(TOPIC_NAME, Consumed.with(stringSerde, ruleConfigSerde)); // 按SCHEME分组 KGroupedStream<String, RuleConfig> groupedKStream = ruleConfigKStream.groupBy( (key, value) -> value.getScheme(), Grouped.with(Serdes.String(), ruleConfigSerde) ); // 聚合为列表,正确配置Materialized KTable<String, List<RuleConfig>> ruleStore = groupedKStream.aggregate( // 初始化空列表 ArrayList::new, // 聚合逻辑:创建新列表避免修改原状态,可选按RULEORDER排序 (key, value, currentList) -> { List<RuleConfig> newList = new ArrayList<>(currentList); newList.add(value); newList.sort(Comparator.comparingInt(RuleConfig::getRuleOrder)); return newList; }, // 配置状态存储和序列化器 Materialized.<String, List<RuleConfig>, KeyValueStore<Bytes, byte[]>>as(RULE_STORE) .withKeySerde(Serdes.String()) .withValueSerde(listSerde) ); builder.build(); // 查询存储 ReadOnlyKeyValueStore<String, List<RuleConfig>> ruleKVStore = kafkaStreams.store( StoreQueryParameters.fromNameAndType(RULE_STORE, QueryableStoreTypes.keyValueStore()) ); List<RuleConfig> mcRules = ruleKVStore.get("MC");
额外注意事项
- 序列化器配置:确保
listSerde是正确的列表序列化器,可通过以下方式创建:Serde<List<RuleConfig>> listSerde = Serdes.ListSerde(List.class, ruleConfigSerde); - 幂等性处理:如果上游可能发送重复记录,需在聚合逻辑中通过
RULEORDER判断是否已存在该规则,避免重复添加。 - 状态刷新时机:查询前需确保Kafka Streams处于
RUNNING状态,等待初始化完成并处理完已有记录。
内容的提问来源于stack exchange,提问作者poonam kumari
相关产品推荐
相关产品推荐

