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

Flink Exactly Once事务超时时仍提交成功的原因咨询

我在使用Flink Exactly Once(EO)模式时,重启Flink作业收到如下警告:

2023-12-30 13:07:44.538 [Co-Flat Map -> Sink: Sink1 (3/8)#0] WARN
o.a.f.s.api.functions.sink.TwoPhaseCommitSinkFunction - Transaction
KafkaTransactionState [transactionalId=ABCD, producerId=1028910,
epoch=3] has been open for 2813524 ms. This is close to or even
exceeding the transaction timeout of 900000 ms.

但随后该事务提交却成功了:

2023-12-30 13:07:44.907 [Co-Flat Map -> Sink: Sink1 (3/8)#0] INFO
o.a.f.s.api.functions.sink.TwoPhaseCommitSinkFunction -
FlinkKafkaProducer 3/8 committed recovered transaction
TransactionHolder{handle=KafkaTransactionState [transactionalId=ABCD,
producerId=1028910, epoch=3], transactionStartTime=1703938851014}

我无法理解为何事务已超过配置的transaction.max.timeout.ms超时时间,提交仍能成功。我查看了TwoPhaseCommitSinkFunction的代码:

private void recoverAndCommitInternal(TransactionHolder<TXN> transactionHolder) {
        try {
            logWarningIfTimeoutAlmostReached(transactionHolder);
            recoverAndCommit(transactionHolder.handle);
        } catch (final Exception e) {
            final long elapsedTime = clock.millis() - transactionHolder.transactionStartTime;
            if (ignoreFailuresAfterTransactionTimeout && elapsedTime > transactionTimeout) {
                LOG.error(
                        "Error while committing transaction {}. "
                                + "Transaction has been open for longer than the transaction timeout ({})."
                                + "Commit will not be attempted again. Data loss might have occurred.",
                        transactionHolder.handle,
                        transactionTimeout,
                        e);
            } else {
                throw e;
            }
        }
    }

    private void logWarningIfTimeoutAlmostReached(TransactionHolder<TXN> transactionHolder) {
        final long elapsedTime = transactionHolder.elapsedTime(clock);
        if (transactionTimeoutWarningRatio >= 0
                && elapsedTime > transactionTimeout * transactionTimeoutWarningRatio) {
            LOG.warn(
                    "Transaction {} has been open for {} ms. "
                            + "This is close to or even exceeding the transaction timeout of {} ms.",
                    transactionHolder.handle,
                    elapsedTime,
                    transactionTimeout);
        }
    }

这个警告难道不意味着提交肯定会失败吗?我是否忽略了什么?甚至作业启动1小时后,旧事务仍能提交?

这是否因为Flink正在恢复事务?Kafka是否允许通过相同的producer id和epoch恢复已中止的消息?相关代码如下:

producer = initTransactionalProducer(transaction.transactionalId, false);
                producer.resumeTransaction(transaction.producerId, transaction.epoch);
                producer.commitTransaction();

如果不是超时,那哪些参数会决定提交失败?

补充说明:transaction.abort.timed.out.transaction.cleanup.interval.ms 使用默认值。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 04:07:05