【发布时间】:2021-07-25 11:19:09
【问题描述】:
我看到 spring-kafka 使用 @RetryableTopic 支持非阻塞重试。我只看到@RetryableTopic 正在与@kafkaListener 一起工作。但我想要我的流聚合上的“重试”。 spring-kafka 如何做到这一点?
下面的代码示例是关于银行交易(流)和账户余额(状态)的。
假设银行交易是这样的:将 10 美元从账户 10001 转移到账户 10002。
我有下面的流代码,使用 reduce 函数从 10001 到 -10 和 +10 到 10002。
余额被物化到状态存储 BALANCE。
如果10001账户余额小于10,则交易不成交。
但是应该重试,因为存款交易可能会在很短的时间内到来。并存入10001后余额>10,则该笔交易完成。
这是我的流豆
@Bean
public KStream<String, BankTransaction> alphaBankKStream(StreamsBuilder streamsBuilder) {
JsonSerde<BankTransaction> valueSerde = new JsonSerde<>(BankTransaction.class);
KStream<String, BankTransaction> stream = streamsBuilder.stream(Topic.TRANSACTION_RAW,
Consumed.with(Serdes.String(), valueSerde));
KStream<String, BankTransaction>[] branches = stream.branch(
(key, value) -> isBalanceEnough(value),
(key, value) -> true /* all other records */
);
branches[0].flatMap((k, v) -> {
List<BankTransactionInternal> txInternals = BankTransactionInternal.splitBankTransaction(v);
List<KeyValue<String, BankTransactionInternal>> result = new LinkedList<>();
result.add(KeyValue.pair(v.getFromAccount(), txInternals.get(0)));
result.add(KeyValue.pair(v.getToAccount(), txInternals.get(1)));
return result;
}).filter((k, v) -> !Constants.EXTERNAL_ACCOUNT.equalsIgnoreCase(k))
.map((k,v) -> KeyValue.pair(k, v.getAmount()))
.groupBy((account, amount) -> account, Grouped.with(Serdes.String(), Serdes.Double()))
.reduce(Double::sum,
Materialized.<String, Double, KeyValueStore<Bytes, byte[]>>as(StateStore.BALANCE).withValueSerde(Serdes.Double()));
return stream;
}
private boolean isBalanceEnough(BankTransaction bankTransaction) {
// read balance from state store BALANCE
return balance >= bankTransaction.amount
}
【问题讨论】: