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)); }
代码说明
- 替换
map为flatMap:flatMap支持将单条输入记录转换为多条输出记录,完美匹配“一个Transaction对应多个AccountDetails生成多条记录”的需求 - 遍历
transaction.getAccounts()中的每个账户,为每个账户生成对应的KeyValue对象,Key按要求拼接格式,Value保留原Transaction对象 - 修正返回值类型:原代码错误返回
List<KStream<...>>,实际应直接返回KStream<...>,符合Spring Cloud Stream KStream binder的单流输入输出规范
内容的提问来源于stack exchange,提问作者Javanectar
相关产品推荐
相关产品推荐

