You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Cloud Kafka:如何基于KStream生成多不同Key的输出记录

Spring Cloud Kafka Binders 实现单输入多输出KStream记录

需求说明

使用Spring Cloud Kafka Binders处理输入的KStream<String, Transaction>数据,为每个AccountDetails生成一条独立的输出记录。

数据结构

@Getter 
@Setter 
public class Transaction {     
  private String paymentRef;     
  private List<AccountDetails> accounts; 
}  
@Getter
@Setter
public class AccountDetails {     
   private String accountId;     
   private String accountHolderName; 
}

输入示例

{
  "paymentRef":"12345", 
  "accounts": [{
    "accountId": "AB12345",
    "accountHolderName": "Dietrich Fritzi"
   }, {
    "accountId": "CD98741",
    "accountHolderName": "Wibke Elisabeth"
   }]
}

期望输出

为输入Transaction中的每个AccountDetails生成一条记录:

  • Key格式为paymentRef-accountId(示例中对应12345-AB12345和12345-CD98741)
  • Value保持原Transaction对象不变

当前问题

现有代码仅能生成一条输出记录,原因是使用findFirst()只取了第一个AccountDetails,未遍历所有账户生成对应记录。

当前错误代码

@Bean
public Function<KStream<String, Transaction>, List<KStream<String, Transaction>>> accountTransaction() {
return transaction -> transaction
        .peek((key, value) -> log.info("Incoming Transaction, Key :: [{}] Value :: [{}]", key, value))
        .map((key, value) -> new KeyValue<>(value.getPaymentRef() + "-" + value.getAccounts().stream().findFirst().map(AccountDetails::getAccountId).orElse("NO_VALUE"), value))
        .peek((key, value) -> log.info("Outgoing Transaction, key :: {} Value :: {}", key, value));
}

正确实现代码

@Bean
public Function<KStream<String, Transaction>, KStream<String, Transaction>> accountTransaction() {
    return transactionStream -> transactionStream
            .peek((key, value) -> log.info("Incoming Transaction, Key :: [{}] Value :: [{}]", key, value))
            // 用flatMap将单条记录展开为多条,每个AccountDetails对应一条输出
            .flatMap((key, transaction) -> 
                transaction.getAccounts().stream()
                        .map(account -> new KeyValue<>(
                                transaction.getPaymentRef() + "-" + account.getAccountId(),
                                transaction
                        ))
                        .collect(Collectors.toList())
            )
            .peek((key, value) -> log.info("Outgoing Transaction, key :: {} Value :: {}", key, value));
}

代码说明

  1. 替换map为flatMap:flatMap支持将单条输入记录转换为多条输出记录,完美匹配“一个Transaction对应多个AccountDetails生成多条记录”的需求
  2. 遍历transaction.getAccounts()中的每个账户,为每个账户生成对应的KeyValue对象,Key按要求拼接格式,Value保留原Transaction对象
  3. 修正返回值类型:原代码错误返回List<KStream<...>>,实际应直接返回KStream<...>,符合Spring Cloud Stream KStream binder的单流输入输出规范

内容的提问来源于stack exchange,提问作者Javanectar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 12:10:01