Flink(1.11.3)对接Kafka(2.1.1) Sink重启后报UnsupportedVersionException求助
结合你使用的Flink 1.11.3和Kafka 2.1.1版本,以及Exactly-Once语义的场景,这些异常本质上都是事务恢复阶段的状态不匹配导致的,下面逐个拆解原因并给出针对性的解决办法:
一、异常原因分析
1. InvalidTxnStateException(事务无效,推测已提交)
你设置的transaction.timeout.ms=30000(30秒)太短了。当JobManager/TaskManager故障重启时,Flink需要时间恢复状态、重新连接Kafka并尝试续上之前的事务,但如果恢复耗时超过30秒,Kafka会自动终止超时的事务并标记为已提交/中止,此时Flink再去操作这个已经失效的事务,就会触发这个异常。
2. InvalidPidMappingException(ProducerID未关联事务ID)
Kafka的事务ID(Transaction ID)和生产者ID(Producer ID)是一一绑定的。当事务超时被Kafka清理后,这种绑定关系会被解除。Flink重启后仍然使用旧的事务ID去请求Producer ID,就会出现“找不到对应映射”的错误。
3. UnsupportedVersionException(版本1下写入非默认ProducerID)
Kafka 2.1.1默认的消息格式版本是v1(对应Kafka 0.11.x),虽然v1支持事务,但Flink在Exactly-Once模式下会自定义生成Producer ID,而旧版本的Kafka在处理这种非默认Producer ID的事务写入时,可能存在兼容性的细节问题——尤其是当事务状态日志(__transaction_state)的版本和Producer的消息格式不匹配时。
二、解决办法
1. 调整事务超时时间(最关键的一步)
把transaction.timeout.ms调大,确保超过Flink任务的最大可能重启恢复时间,同时要保证Kafka Broker的transaction.max.timeout.ms(默认900000,15分钟)大于等于这个值,否则Broker会拒绝Producer的事务请求。
修改你的Producer配置:
props.setProperty("transaction.timeout.ms", "900000"); // 15分钟
同时检查Kafka Broker的server.properties,确保:
transaction.max.timeout.ms=900000
2. 确保Kafka事务日志的高可用性
Kafka的事务状态都存在__transaction_state主题中,如果这个主题的副本数或ISR(同步副本)不足,会导致事务状态无法正确持久化,进而在恢复时出现异常。生产环境建议配置:
transaction.state.log.replication.factor=3 transaction.state.log.min.isr=2
如果这个主题还没自动创建,可以手动创建,确保分区数为默认的5,且所有分区状态正常。
3. 显式指定Kafka消息格式版本
手动指定Producer的消息格式版本为Kafka 2.1.1对应的2.1,避免版本自动协商带来的兼容性问题:
props.setProperty("message.format.version", "2.1"); props.setProperty("inter.broker.protocol.version", "2.1"); // 和Broker配置保持一致
同时确保Producer的序列化器配置正确(你的代码里已经用了ByteArraySerializer,这个没问题)。
4. 优化Flink重启策略
避免快速频繁重启导致的事务恢复冲突,在flink-conf.yaml中配置固定延迟重启:
restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 3 restart-strategy.fixed-delay.delay: 10s
给Kafka足够的时间清理失效事务,再进行恢复尝试。
5. 检查Flink状态后端配置
确保Flink使用持久化的状态后端(比如RocksDB),这样重启时能正确恢复事务的上下文信息,避免因为状态丢失导致的事务恢复混乱。
三、额外注意点
如果以上配置调整后仍然出现问题,可能是Flink 1.11.x版本的Kafka Producer存在事务恢复的小bug——毕竟1.11是比较老的版本了。如果业务允许,升级到Flink 1.12+版本会有更好的事务兼容性,但如果暂时无法升级,通过上面的配置调整基本可以覆盖大部分场景。
内容的提问来源于stack exchange,提问作者jinsunism

