惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

博客园_首页
H
Help Net Security
N
Netflix TechBlog - Medium
Apple Machine Learning Research
Apple Machine Learning Research
P
Proofpoint News Feed
A
About on SuperTechFans
V
V2EX
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
宝玉的分享
宝玉的分享
aimingoo的专栏
aimingoo的专栏
F
Fortinet All Blogs
博客园 - 【当耐特】
Microsoft Security Blog
Microsoft Security Blog
Martin Fowler
Martin Fowler
I
InfoQ
Google DeepMind News
Google DeepMind News
人人都是产品经理
人人都是产品经理
Engineering at Meta
Engineering at Meta
腾讯CDC
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
B
Blog RSS Feed
U
Unit 42
The Cloudflare Blog
Y
Y Combinator Blog

Spaceack's Blog

使用简单算法两小时实现猎杀乌姆帕斯(Hunt the Wumpus)Python小游戏 搭建etcd分布式kv数据库集群并使用python操作 使用docker模拟双机环境测试mysql双机互为主备 天池大数据竞赛 Spaceack带你利用Pandas,趋势图与桑基图分析美国选民候选人喜好度 Python利用matplotlib万花尺画月饼 情人节遥寄 靠近爱人 牵着你的手 七夕收到的礼物
使用Mysql实现消息队列
Spaceack · 2021-09-30 · via Spaceack's Blog

/20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/1.gif

实现起来就是 消息状态版本号 字段。

/20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/2.png

更新时用 版本号 做乐观锁。操作逻辑就是个状态机。

UPDATE mq SET mq.status=new_status mq.version = mq.version + 1 WHERE mq.version = old_version

实现

mysql mq 表结构设计

1
2
3
4
5
6
7
CREATE TABLE `mq` (
  `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
  `msg` varchar(1024) DEFAULT NULL,
  `status` varchar(100) DEFAULT 'ready',
  `version` bigint(20) unsigned DEFAULT 0,
  PRIMARY KEY (`id`)
) ENGINE=InnoDB AUTO_INCREMENT=10 DEFAULT CHARSET=utf8mb4;

生产者

测试,向队列插入两条消息。

1
2
insert into mq(msg) values('第一条消息 hello world! 1024');
insert into mq(msg) values('第二条消息 欢迎访问 https://spaceack.com');

/20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/3.png

消费者

  • 获取队列中的消息, 此时不会改变队列(mq表)中的数据。

    1
    
    select * from mq where status='ready' limit 1;
    

    /20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/4.png

  • 消息确认

    1
    
    update mq set mq.status = 'ack', mq.version = mq.version + 1 WHERE mq.version = {query_version} and id = {query_id}
    

    确认后的状态:

    /20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/5.png

再次获取数据仅能获取第二条数据。

/20210929-%E4%BD%BF%E7%94%A8mysql%E5%AE%9E%E7%8E%B0%E6%B6%88%E6%81%AF%E9%98%9F%E5%88%97/6.png


这样的一个好处就是消息都是可见的。 可以直接查数据库。

不用加悲观锁,因为update成功后就带锁。其它同版本的更新都会失败。

还有个好处就是减少组件依赖。 简单的服务数据库就能搞定。 不用再起个rabbitmq服务啥的。 节约运维成本。

版本号 另一个小作用:InnoDB如果更新语句没有改变任何字段值时,影响行数会返回0.那么是没找到记录还是没改变值的0呢?前者是bug, 后者是正常情况但区分不了。加版本号后每次修改+1,影响行数一定不为0。


隐藏福利:生产者自动批量测试脚本:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
import random
import subprocess

def runcmd(command):
    ret = subprocess.run(command,shell=True,stdout=subprocess.PIPE,stderr=subprocess.PIPE,encoding="utf-8",timeout=1)
    if ret.returncode == 0:
        # print("success:",ret, ret.stdout)
        return ret.stdout
    else:
        print("error:",ret)


def producer():
    sql = "\"insert into mq(msg) values('%s');\"" % (str(random.randint(1,99999)))
    cmd = """mysql -uroot -ppassword test -e %s """ % (sql)
    runcmd(cmd)