ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka事务消息脏读事故:read_committed与read_uncommitted隔离级别解析

Kafka事务消息脏读事故:read_committed与read_uncommitted隔离级别解析 1. 事故现场消息队列里凭空多出来的“鬼消息”1.1 现象日志一切正常数据却对不上账复盘这事之前先描述一下事故当天的情况。我们有个订单状态同步链路A系统负责接收支付回调把订单状态发到KafkaB系统负责消费这些消息落库后再触发下游的发货、开票等动作。某天运维报警B系统的数据库里出现了一大批“超前”的状态记录——订单状态已经被改成了“已付款”但A系统里对应的订单根本没有完成支付甚至有一批在几分钟后被业务侧标记成了“已取消”。这类问题最磨人的地方在于你去看日志发现链路里每个环节都正常。A系统发了消息日志显示send成功B系统消费了消息日志显示处理完成offset正常提交Kafka监控面板上消费堆积为0延迟也不高。数据库里的记录看起来也很干净没有异常堆栈。可业务对账就是不平多出来好多“凭空出现”的已付款订单。第一天我们基本是在业务代码里找原因。怀疑过消息重复、怀疑过状态机更新顺序写错、怀疑过下游代码有位运算bug反复看B系统的消费逻辑发现它很简单拿到订单状态消息做幂等校验然后update数据库。没有复杂分支也没有可能写错状态的逻辑。真正让我把视线转向Kafka的是一个细节这些“鬼消息”对应的订单恰好都在A系统最近改过的一个事务发送模块里。1.2 时间线复盘上下游都认为自己没错把时间线拉出来之后事情开始变得清晰。我用Kafka的消费时间戳和生产端业务日志做了对齐发现一个很刺眼的差值A系统在14:02:31执行beginTransactionA系统在14:02:35把订单状态消息send到KafkaA系统直到14:03:12才执行commitTransactionB系统的消费时间却是14:02:34比事务提交早了38秒这个38秒就是问题核心。Kafka在设计上事务型生产者写入的消息在事务提交前对于常规消费链路是不应该可见的。可我们的事实是事务还在进行中消费者已经读到了这条消息并且正常消费、正常落库了。把上下游各自视角列出来看两边都没有错。A系统认为“我事务提交成功了消息肯定没问题”B系统认为“我按正常逻辑消费没报异常也没堆积”但把两个视角放在一起中间的可见性规则出现了裂缝。所有排查Kafka问题的人都要记住这类故障不会直接报错它更像一种“静默的约定破坏”——单个节点永远正常整体数据却对不上。2. 生产端为什么“开了事务却没开彻底”2.1 事务型生产者的正确配置全貌先梳理一遍Kafka事务的正常姿势。它和数据库事务差很远不是发一条BEGIN语句就完事。Kafka事务需要生产者在启动时明确声明一个全局唯一的transactional.id然后走一套“初始化、开启、写入、提交/回滚”的完整流程。基础配置长这样Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka1:9092,kafka2:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, order-txn-producer-001); props.put(ProducerConfig.ACKS_CONFIG, all); KafkaProducerString, String producer new KafkaProducer(props); producer.initTransactions();发送逻辑的样板是try { producer.beginTransaction(); producer.send(new ProducerRecord(order-status, orderId, PAID)); producer.send(new ProducerRecord(order-status-log, orderId, paid at 14:02:35)); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); throw e; }这里有两个关键点容易踩坑。第一transactional.id必须稳定且唯一它不仅是身份标识还承载着事务恢复的职责。如果同一个transactional.id被两个生产者实例同时使用Kafka会触发ProducerFencedException把后启动的实例“围栏”掉。第二acks要配置成all事务机制虽然强制开启了幂等但只有acksall才能确保副本同步避免主副本切换时丢消息。Kafka事务本身会在底层强制设置enable.idempotencetrue这是与旧版本最大差别之一。如果你在代码里试图把幂等关掉初始化阶段就会直接报错所以这两处的配置项要看成一体。2.2 Spring Kafka里的事务边界它和数据库事务不是一回事我们线上用的是Spring Boot所以绕不开Spring Kafka的事务写法。很多人在Spring里写的“事务”其实只是Transactional注解加在Service方法上但Kafka的KafkaTemplate并不会自动参与Spring的数据库事务除非你做了明确的绑定。Spring Kafka提供了一套机制把KafkaTransactionManager配置为PlatformTransactionManager然后Transactional注解才能同时管住数据库事务和Kafka消息发送。如果只是简单地在KafkaTemplate.send外面套一个Transactional而没有配置对应的KafkaTransactionManager那这个注解对Kafka来说相当于不存在事务不会被真正开启。举一个我们实际踩过的简化版代码Transactional public void processPayment(Order order) { orderMapper.updateStatus(order.getId(), PAID); kafkaTemplate.send(order-status, order.getId(), PAID); }如果项目里没有配置KafkaTransactionManager这段代码里的数据库更新会走Spring的数据库事务但kafkaTemplate.send仍然是普通发送根本不会进入Kafka事务流程。等到后来有人把这段代码补上了KafkaTransactionManager消息发送才真正进入事务但新的问题也随之而来消费端完全不知情。所以当你在排查“生产端明明开了事务但消费者还是读到了怪消息”时先确认一件事生产端到底是真的开了Kafka事务还是只是你以为开了。如果是Spring配置漏了KafkaTransactionManager那生产端根本没有事务消费者读到什么都是正常的。但这次我们不是这种情况因为从broker侧能看到事务控制消息说明事务确实开了。2.3 事务超时与协调器另一个看不见的坑Kafka事务还有一个容易被忽略的配置transaction.timeout.ms默认是60000ms。这个参数决定事务从开启到提交允许的最大时间。如果超过这个时间还没提交事务协调器会主动中止事务。这会导致一个很有意思的现象生产端代码里明明没调用abortTransaction事务却被broker强制标记为aborted。消费者如果跑在read_uncommitted级别照样会消费到这批消息。更坑的是生产端日志可能只显示“send成功”压根不提示事务被超时回滚因为发送操作本身已经完成了。我们这次事故里没有直接触发超时但排查时我特意去看了transaction.state相关指标确认没有超时中止的记录。如果你在处理类似问题这一步不能省。Kafka的事务协调器会记录每个事务的状态机转换通过JMX指标或者kafka-transactions日志可以查看到完整的提交、中止记录这些信息能帮你区分“消息本身被abort”和“消费者错误地消费了已提交消息”两种情况。2.4 消费端隔离级别整条链路最容易被忽视的环节回到事故本身。生产端配置没问题事务也确实开了为什么消费者还能提前读到原因就在消费者端的isolation.level。Kafka消费者有两个隔离级别read_uncommitted默认值。所有消息都能看包括未提交事务里的消息以及最终被abort回滚的消息。read_committed只能看到已提交事务的消息以及非事务型生产者发送的普通消息未提交和已回滚的事务消息会被过滤。我们B系统的代码里isolation.level没有做过任何配置。这个配置项在Java Consumer里的默认值就是read_uncommitted所以消费者根本没有理会事务状态直接就把消息消费了。等到生产端事务回滚这批消息已经在下游数据库里产生了不可逆影响Kafka本身不会帮你“撤销”消费者落库的数据。这里要特别提醒用Spring Boot的同学spring.kafka.consumer.properties.isolation.level如果不显式配置默认也是read_uncommitted。KafkaListener不会帮你自动切换成read_committed。很多人觉得“Kafka事务嘛生产端配好就行了”消费端这个配置恰恰是整条事务链路的另一半。提示当生产端引入Kafka事务后消费端必须显式设置isolation.levelread_committed并且这个配置要覆盖所有消费该Topic的消费者组一个都不能漏。3. read_committed与read_uncommitted的底层逻辑3.1 LSO与事务控制消息read_committed到底在过滤什么不理解底层机制很容易把read_committed当成一个“魔法开关”。实际上它的过滤逻辑完全基于offset和事务控制消息并不复杂。Kafka的事务消息不会单独存放在特殊区域它和普通消息一样写入分区只是消息流里掺杂了控制消息control records。当生产端开启事务时会在每个涉及的分区写入一条事务开始的控制消息提交时写入提交控制消息回滚时写入中止控制消息。这些控制消息对read_uncommitted消费者是完全透明的它只按offset顺序读数据消息控制消息直接跳过。read_committed消费者则维护了一个叫做LSOLast Stable Offset最后稳定位移的指针。LSO代表当前分区中“已稳定”消息的边界只有小于LSO的消息才是可以安全消费的。消费者会持续跟踪事务状态如果遇到了一个未完成的事务LSO就停在这个事务开始的控制消息位置不再向下推进直到收到事务提交或中止的控制消息才会把LSO推进到事务结束位置。这个机制带来的直接后果是read_committed消费者可能因为某个长时间未提交的事务而卡住看到消费延迟升高但这是为了消息可见性付出的代价。read_uncommitted消费者没有这个限制它永远不被事务状态阻塞所以读取实时性更高代价就是会读到“脏数据”。3.2 一个让大多数人意外的特性abort消息不会消失我发现很多Kafka开发者有一个根深蒂固的误解以为事务回滚后消息会被从Kafka里删除。实际上Kafka几乎不会因为事务回滚而删除任何消息。消息一旦写入分区就成为segment文件里的一段不可变记录。事务abort只是在控制消息层面标记“这个事务中止了”涉及的数据消息仍然原封不动留在分区里。read_committed消费者看到中止控制消息后会从逻辑上跳过这些数据消息read_uncommitted消费者则完全不会过滤直接把消息读出来交给业务代码。所以在我们的场景里A系统有一批订单因为校验失败走入了abortTransaction但B系统读到了这些被回滚的消息把订单状态更新成了“已付款”。从业务角度这比“提前消费已提交消息”严重得多——前者至少最终状态是一致的只是时间提前后者根本是拿了一个“被废弃”的数据去更新真实业务。这里可以打个比方Kafka分区就像仓库货架事务abort的消息像一批贴了“作废”标签的箱子。read_committed消费者只搬“入库单”盖章的箱子read_uncommitted消费者不管标签见箱子就搬。仓库管理员不会因为标签作废就把箱子烧掉它还是躺在货架上只是等你决定要不要搬而已。3.3 混合消息与跨分区问题read_committed不是万能的顺带说一个边界情况。read_committed只对“事务型生产者写入的消息”做过滤对非事务型生产者写入的普通消息它不做任何事务判断直接视为可见。所以如果同一个Topic里既有事务消息又有普通生产者发送的消息消费端不会把普通消息误过滤掉这点可以放心。但跨分区消费时会有另一个问题Kafka事务可能跨多个分区消费者在某个分区上看到事务开始的控制消息在另一个分区上可能事务还没开始。LSO是逐分区维护的read_committed消费者读取每个分区时会按照“该分区最后稳定位移”来限制读取范围。这意味着即使消息已经提交不同分区的消费者看到它的时间点也可能略有差异。对于大部分业务场景这个差异可以忽略但如果你的业务对“消息可见时刻”有强一致要求需要把视角放大到整个事务的协调状态而不是只看单个分区。这次排查时我就发现B系统消费的Topic有6个分区A系统发送的消息在4个分区上分布不均但好在每一条被污染的订单都能通过消息key定位到具体分区这才让我们能精确统计影响范围而不是全量重放。4. 完整排查链路从一切正常到锁定根因4.1 第一步先用消费组工具确认“肚子里的数据”排查Kafka问题我习惯从最外层开始。第一个动作是查消费者组的offset情况bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \ --describe --group order-status-consumer输出结果如表所示脱敏后的关键字段分区当前offsetlog-end-offsetlag01043210432019877987702108351083503105011050104100221002205986698660这份结果非常有欺骗性每个分区lag都是0说明消费者已经把Topic里所有消息消费完了没有任何积压。当时团队里有人看到这个就说“消费端没问题”差点就把方向带偏了。lag为0只能说明消费进度不能说明消费内容的正确性。消费者消费了一条aborted消息它对Kafka而言也是“消费完了”offset正常提交lag自然为0。所以这一步的真正价值是确认“消费者确实消费了所有消息”再结合时间线下一步就要去原始分区里看消息到底长什么样。4.2 第二步翻原始分区找到abort标记第二步我直接用控制台消费者脚本去读原始消息不经过业务代码不带任何消费组从最早的offset开始扫bin/kafka-console-consumer.sh --bootstrap-server kafka1:9092 \ --topic order-status \ --property print.keytrue \ --property print.valuetrue \ --property print.partitiontrue \ --property print.offsettrue \ --from-beginning \ --max-messages 5000控制台消费者默认使用read_uncommitted所以它能显示出所有消息包括可能被事务过滤掉的“鬼消息”。我把时间范围卡在14:02到14:03这个窗口重点看分区0的offset 10430到10490区间。这段消息里数据消息夹杂着一些特殊记录值部分为空、key也是空但它们的offset存在。这些就是事务控制消息用普通业务消费者读出来往往会被跳过但用控制台消费者配print.offset能看到它们的痕迹。再对照A系统生产端的事务日志我发现一个规律offset 10435和10436之间有事务开始的控制消息offset 10490和10491之间是事务中止的控制消息而10436到10490之间的数据消息正是B系统数据库里多出来的“已付款”订单。这些消息在read_uncommitted视角下被完整消费了但在事务层面它们是处于aborted状态的。到这里根因基本浮出水面了。4.3 第三步写一个最小验证消费者做对照实验证据链还不够稳我用Java写了一个几十行的最小验证消费者专门验证两个隔离级别的差异。Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka1:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, verify-group-read-uncommitted); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, read_uncommitted); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest);先用read_uncommitted消费统计消费到的消息总数然后换一个组名把ISOLATION_LEVEL_CONFIG改成read_committed再消费一遍统计总数。两边对比差值正好是18条全部落在aborted事务区间。这一步做完整个排查链路就闭环了。不是业务代码的问题不是生产端发送失败的问题不是Kafka集群broker的问题是消费端隔离级别配置缺失导致的事务消息可见性泄漏。这个结论不再停留在“我觉得是”的层面而是可以用两个消费组的数据明确复现出来的。4.4 为什么上游一直没发现事务正常率还在高位整个排查过程中A系统的同学一直很委屈他们的事务成功率明明超过98%告警阈值设在95%根本没触发。后来看详细指标那1.7%的abort事务都集中在某个外部接口偶发超时的时间窗口平时一个月也就出现几次不会引起注意。但就是这么低的abort比例因为Kafka消费者错误地读取了aborted消息造成了下游数据库的18条脏记录。这个放大效应是很多团队想不到的数据库事务回滚数据对任何人不可见影响很小Kafka事务回滚只要消费端不是read_committed脏数据照样流向下游影响被放大了几十倍。所以遇到Kafka事务类问题别只盯着“事务成功率”这种生产端指标消耗端的可见性指标同样重要。如果生产端abort数量大于0就应该意识到所有消费该Topic的消费者组都有潜在的脏读风险。5. 修复方案、验证与复盘中的关键细节5.1 线上修复改配置只是第一步修复方案本身不复杂消费端加上隔离级别配置就行。用Spring Boot的话改application.ymlspring: kafka: consumer: properties: isolation.level: read_committed如果用原生API则在消费者Properties里加一行props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, read_committed);加了配置之后需要滚动重启消费服务。注意一个容易搞错的点修改isolation.level不会重置消费者的offset它只是改变了后续消费的可见性规则。原本已经消费过的offset不会回退Kafka会从当前已提交的offset继续消费。所以“改了配置重启就好”这个说法不完整因为已经被污染的下游数据不会自动恢复。重启后我们用一个验证消费者重新消费了一遍故障时间窗口内的消息确认aborted消息不再出现在read_committed模式下这一步验证了配置生效。然后才进入数据补偿阶段。5.2 数据补偿比改配置更麻烦的活补偿这18条脏数据比改配置麻烦得多。因为错误的“已付款”状态已经落库后续可能触发了发货流程、发票流程、物流通知等联动动作直接UPDATE数据库是一种危险操作。我们最后走的方案是从原始消息和下游数据库记录里提取出受影响的订单号清单。对照A系统的最终订单状态确定每条订单应回退到“待支付”还是“已取消”。通过业务对冲接口逐个回退状态并记录审计日志而不是绕过业务逻辑直接改库。对于已经触发联动流程的订单由业务团队单独人工介入处理防止二次污染。这个补偿过程花了两天时间比根因修复本身的时间长得多。它说明了一个道理Kafka事务类故障的修复成本大部分不在改配置而在清洗已经被污染的下游数据。如果业务链路再长一点比如状态又触发了库存扣减、优惠券核销那补偿的复杂度会指数级上升。5.3 回归验证不能只看日志要看消息流修完之后我设计了一个回归用例避免以后再犯同样的错。场景分三种正常事务提交的消息消费端必须能消费到。正常事务abort的消息消费端绝不能消费到。非事务型生产者发送的普通消息消费端必须能消费到。我写了一小段验证程序一个事务型生产者先发送一条消息但不提交sleep 10秒后abort同时跑一个read_committed消费者订阅同一个Topic。预期结果是消费者在事务abort前后都不该看到这条消息。实测通过。再跑一个用例发送并提交事务消费者正常收到消息。也通过。回归用例通过后B系统才能重新进入正常业务流。同时我保留了故障时间窗口的原始消息dump放到日志平台上方便后续其他团队对账时回溯。提示read_committed模式下如果某个事务长时间不提交消费者会阻塞在LSO位置表现为消费延迟上升。这不是故障而是隔离级别带来的正常行为。排查延迟时要先确认是否有人在生产端开了长事务。6. 给团队定下的几条硬规矩6.1 代码评审检查项事务型生产者必须检查消费端这次复盘后我把Kafka事务相关检查项正式写进了代码评审规范只要有涉及Kafka事务的改动评审时必须逐项确认生产端是否配置了transactional.id且全环境唯一。生产端是否明确配置了acksall。事务的commit和abort是否覆盖所有分支尤其是网络超时、外部接口异常。所有消费该Topic的消费者组是否都显式设置了isolation.levelread_committed。是否有非事务型生产者向同一Topic发送消息是否会破坏事务语义。是否存在跨集群、跨机房场景下的事务协调器连接问题。这些检查项没有一条是复杂的但每一条都能拦住一次实实在在的事故。尤其是第四项只看生产端不看消费端是Kafka事务项目最常见的翻车姿势。6.2 监控与快速应急事务可见性要做到可观测除了评审监控也要跟上。我们增加了三个维度第一生产端事务状态监控。按分钟统计beginTransaction、commitTransaction、abortTransaction三种状态的次数abort占比超过1%即告警。这个阈值可以根据业务调整但一定要有。第二消费端隔离级别巡检。写一个定时任务扫描所有消费组的配置如果发现某个消费者组没有配置read_committed但订阅了使用了Kafka事务的生产者Topic自动提示风险。这个东西实现起来不难但对“多消费者组共用Topic”的场景特别有价值。第三故障演练。每季度安排一次Kafka事务abort演练构造一个小规模事务主动abort然后观察所有相关消费者组是否都能正确过滤。成本很低十几分钟就能完成但能提前把“消费端没配隔离级别”这类问题暴露在演练环境而不是等业务对账时再发现。6.3 一些写给后来者的碎碎念如果你正在看这篇文章并且你们正准备用Kafka事务来解决“本地消息和Kafka发送不一致”的问题我的建议是别急着写代码先把消费端隔离级别这条配套规则同步给所有下游团队。Kafka事务本质上是一套“端到端一致性协议”它要求生产端和消费端共同遵守约定。生产端开了事务消费端就必须用read_committed来配合否则事务的“一致性”就是一句空话。这次事故里我们只改了生产端消费端没跟上结果一条错误状态订单从Kafka漏到了下游数据库最终靠人工补偿才恢复。我个人排查这类问题的最大体会是Kafka事务类故障都不是“报错型”故障而是“静默型”故障。日志里没有异常监控面板上看不到堆积只有业务对账时才发现数据对不上。所以排查时一定要跳出“看日志、看报错”的惯性拉时间线、查原始消息、做对照实验用证据链把故障钉死。最后再分享一个实用小技巧排查前先看一眼所有消费该Topic的消费者组配置isolation.level是不是默认值如果是默认值先跑一个简单的对照实验往往五分钟就能破案。
返回列表