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

