【问题标题】:spring-kafka how to set Retry on the stream beanspring-kafka 如何在流 bean 上设置 Retry
【发布时间】: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
    }

【问题讨论】:

    标签: spring-kafka spring-retry


    【解决方案1】:

    KStream 不在 Spring for Apache Kafka 的范围内; spring 只涉及设置拓扑;一旦设置好,您就可以直接使用 kafka-streams 了。所有功能都由您设置的拓扑提供。

    @RetrybleTopic 功能仅适用于 @KafkaListener(或者更具体地说是 kafka lister 容器)。

    【讨论】:

      猜你喜欢
      • 2021-09-28
      • 2020-09-27
      • 1970-01-01
      • 2023-03-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-03
      相关资源
      最近更新 更多