ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

单应用下RabbitMQ如何保证线程安全,及多应用下抢数据问题

2026/8/6 20:00:29 拓冰建站 浏览量
单应用下RabbitMQ如何保证线程安全,及多应用下抢数据问题 消费RabbitMQ时的注意事项,如何禁止大量的消息涌到Consumerï¼Œä¿è¯çº¿ç¨‹å®‰å ¨ï¼šæŒ‰ç §å®˜ç½‘æä¾›çš„è®¢é˜ åž‹å†™æ³•ï¼ˆ Retrieving Messages By Subscription (push API) ) 我发现,RabbitMQæœåŠ¡å™¨ä¼šåœ¨çŸ­æ—¶é—´å† å‘é€å¤§é‡çš„æ¶ˆæ¯ç»™Consumerï¼Œç„¶åŽï¼Œå¦‚æžœä½ æ²¡æœ‰æ¥å¾—åŠAck的话,那么服务端会积压大量的UnAcked消息,而Consumerå¦‚æžœæ¥ä¸æ€¥å¤„ç†ä¹Ÿä¼šå¤„äºŽå‡æ­»ï¼ˆä¹Ÿå¯èƒ½å¼•èµ·ç¨‹åºå´©æºƒï¼‰ã€‚ä» æœ‰ä¸¤ä¸ªChannel,结果积压了大量的UnAcked消息。这明显是与我们的目的不一致,我们不能保证Consumer一 定会及时快速的处理消息。所以这种方式带来的后果就是Consmer崩溃后,UnAcked消息又ReQueue,这肯定会消耗MQ的宝贵资源。我试图在官网上找到一种方法,让每条消息明确的Ack后再接受下一条。但是没有。好在在 gitbooks.io/rabbitmq-quick/ 这儿找到了,通过设置Channel的QOS即可var channel Connect.CreateModel();channel.BasicQos(0,1,false); //RabbitMQ客户端接受消息最大数量设置后的结果:在开启4个Consumerçš„æƒ å†µä¸‹ï¼Œæ¯æ¡æ¶ˆæ¯å¤„ç†è¦è€—æ—¶2秒。然后问题解决了。Unacked的消息只有4个。让每条消息明确的AckåŽå†æŽ¥å—ä¸‹ä¸€æ¡ï¼Œé€šè¿‡è¿™ç§æ–¹å¼ä¸ä» å¯ä»¥è§£å†³ç§¯åŽ‹å¤§é‡æ¶ˆæ¯çš„é—®é¢˜ï¼Œè¿˜å¯ä»¥é˜²æ­¢å¤šæ¡æ¶ˆæ¯åŒæ—¶æ¶Œå ¥è€Œå¯¼è‡´çš„å¤šçº¿ç¨‹ä¸å®‰å ¨é—®é¢˜ã€‚å¤šåº”ç”¨ï¼šå½“ç„¶åœ¨å•åº”ç”¨ä¸‹ä¸Šé¢çš„æ–¹æ¡ˆå·²ç»å¯ä»¥è§£å†³å¤§éƒ¨åˆ†é—®é¢˜ï¼Œä½†æ˜¯åœ¨å¤šåº”ç”¨ä¸‹ä¼¼ä¹Žä¸å¤ªç†æƒ³ã€‚åˆšå¼€å§‹æ˜¯å•åº”ç”¨æœåŠ¡ï¼Œä½†æ˜¯åŽæ¥ç”¨æˆ·é‡å¹¶å‘é‡æé«˜äº†ï¼Œæ¸æ¸åœ°å•åº”ç”¨æœåŠ¡æ‰›ä¸ä½äº†ï¼Œå°±é€šè¿‡Nginxè´Ÿè½½å‡è¡¡åˆ†æ‹ åˆ°å¦å¤–ä¸€å°æœåŠ¡æœºï¼Œæ­¤æ—¶å°±æ˜¯åŒåº”ç”¨äº†ã€‚è¿™æ ·ä¸€æ¥æœ‰ä¸¤ä¸ªåº”ç”¨åŒæ—¶ç›‘å¬åŒä¸€ä¸ªæ¶ˆæ¯é˜Ÿåˆ—ï¼Œæœ‰ä¸¤ä¸ªæ¶ˆè´¹è€ ã€‚æ¶ˆæ¯é˜Ÿåˆ—çš„å¤„ç†æœºåˆ¶ï¼šä¸€ä¸ªæ¶ˆæ¯é˜Ÿåˆ—é‡Œé¢çš„ä¸€æ¡æ¶ˆæ¯åªä¼šä» ä¸”å‘é€åˆ°ä¸€ä¸ªæ¶ˆè´¹è€ é‡Œé¢ï¼Œç­‰å¾ è¯¥æ¶ˆè´¹è€ Ackï¼Œä¸€æ¡æ¶ˆæ¯å°±å®Œæˆå®ƒçš„ä½¿å‘½äº†ã€‚æœ¬ä»¥ä¸ºè¿™æ ·å­å¤šåº”ç”¨ä¸‹ä¹Ÿä¸ä¼šå‡ºçŽ°ä»€ä¹ˆé—®é¢˜ï¼Œä½†æ˜¯é—®é¢˜å´å‡ºåœ¨åˆ«çš„åœ°æ–¹äº†ã€‚ï¼ˆæ³¨ï¼šåº”ç”¨A与应用B同库)当第一条消息分发到 应用A,第二条消息分发到 应用Bï¼Œä»–ä»¬å‡ ä¹Žæ˜¯åŒæ—¶è¿›è¡Œï¼Œæ¶ˆæ¯A在应用A处理事务,消息B在应用Bå¤„ç†äº‹åŠ¡ï¼Œä»–ä»¬éƒ½éœ€è¦èŽ·å–æ•°æ®åº“ä¸­çš„æŸæ¡æ•°æ®å¹¶ä¸”æ›´æ–°å ¶çŠ¶æ€ï¼Œç”±äºŽA没处理完数据状态未改,Bå› æ­¤ä¹ŸèŽ·å¾—äº†ç›¸åŒçš„ä¸€æ¡æ•°æ®ï¼Œå¯¼è‡´A,B抢了同一条数据做处理了。由于应用A,应用Béƒ½æ˜¯ç‹¬ç«‹çš„æœåŠ¡ï¼Œæ‰€ä»¥å•åº”ç”¨ä¸‹çš„åœ¨ä»£ç é‡Œé¢åŠ åŒæ­¥é”è¿™äº›å¯¹ä»–ä»¬æ¯«æ— ä½œç”¨ã€‚A,B同库,试着利用数据库本身的读写锁机制,进一步优化在insert之前在做一次判断,如果更新成功则insert,否则表明该数据已经被更新了,放弃insertã€‚åŠ äº†è¿™ç§åˆ¤æ–­æ–¹å¼ç¡®å®žå‡å°‘äº†å¾ˆå¤šæŠ¢æ•°æ®é—®é¢˜ï¼Œä½†å¹¶ä¸èƒ½å®Œå ¨è§£å†³è¿™ä¸ªé—®é¢˜ï¼Œæ¯•ç«Ÿæ•°æ®åº“é”åªæ˜¯åœ¨è¯»å†™çš„æ—¶å€™å¯¹ä¸€æ¡æ•°æ®çš„åŠ é”ï¼Œæ“ä½œå®Œäº†é‡Šæ”¾ï¼Œä½†ç”±äºŽåŒæ­¥æ€§å¤ªé«˜ï¼Œå¹¶ä¸èƒ½è§£å†³çŽ°åœ¨çš„é—®é¢˜ã€‚æœ€åŽè€ƒè™‘åˆ°ä¸¤ä¸ªåº”ç”¨éƒ½æ˜¯å®Œå ¨ä¸€æ ·çš„ï¼Œæ‰€ä»¥å°±å¹²æŽ‰ä¸€ä¸ªåº”ç”¨çš„æ¶ˆæ¯æ¶ˆè´¹è€ ï¼Œåªä¿ç•™ä¸€ä¸ªæ¶ˆè´¹è€ ï¼Œè¿™æ ·å°±å¯ä»¥å›žå½’åˆ°å•åº”ç”¨çš„æƒ æ™¯äº†ã€‚æ€»ç»“ï¼šè¿™ä¹Ÿæ˜¯ä¸ºä»€ä¹ˆåœ¨åˆ†å¸ƒå¼ï¼Œå¾®æœåŠ¡ä¸‹åˆ†å¸ƒå¼äº‹åŠ¡çš„å¿ è¦æ€§å’Œé‡è¦æ€§ï¼Œç›®å‰åˆ†å¸ƒå¼äº‹åŠ¡ä¸»è¦é€šè¿‡ MQ事件表、业务补偿、TCCã€å¯¹è´¦ç­‰æ–¹å¼å®žçŽ°ï¼Œä½†éƒ½ä¸å¥½åšï¼Œæ‰€ä»¥å°½é‡é¿å ã€‚åœ¨ä½¿ç”¨åˆ†å¸ƒå¼ï¼Œå¾®æœåŠ¡å¸¦æ¥çš„æ–¹ä¾¿åŒæ—¶ï¼Œä¹Ÿå¾—ä¸ºäº‹åŠ¡çš„å››ä¸ªç‰¹æ€§ï¼ˆåŽŸå­æ€§ï¼Œä¸€è‡´æ€§ï¼Œéš”ç¦»æ€§ï¼ŒæŒä¹ æ€§ï¼‰ä»˜å‡ºä»£ä»·ã€‚å ³æ³¨å ¬ä¼—å·ï¼Œåˆ†äº«å¹²è´§ï¼Œè®¨è®ºæŠ€æœ¯作者molashaonian