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

推荐订阅源

V
Visual Studio Blog
Engineering at Meta
Engineering at Meta
月光博客
月光博客
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
T
Tailwind CSS Blog
博客园 - Franky
The GitHub Blog
The GitHub Blog
大猫的无限游戏
大猫的无限游戏
The Cloudflare Blog
B
Blog RSS Feed
云风的 BLOG
云风的 BLOG
小众软件
小众软件
罗磊的独立博客
Microsoft Azure Blog
Microsoft Azure Blog
I
InfoQ
美团技术团队
H
Hackread – Cybersecurity News, Data Breaches, AI and More
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
V
V2EX
C
Check Point Blog
WordPress大学
WordPress大学
博客园 - 【当耐特】
博客园 - 司徒正美
D
Docker

Jiajun的技术笔记

你好,2026! TiDB 源码阅读(六):TiDB Coprocessor 源码解析 性能优化的核心思想 TiDB 源码阅读(五):索引 TiDB 源码阅读(四):AST、逻辑计划、物理计划 CockroachDB Serverless Architecture podman 无故退出 Cursor Control-L (CTRL-L) Keyboard Shortcuts in Terminal Replace docker with podman Using xmonad with xfce4 A RC script for freebsd frpc 自己动手写一个k8s controller AI 会取代你的(编程)岗位吗? 自建DERP服务器提升Tailscale连接速度(使用Nginx转发) 自动升级Docker容器 再读《程序员修炼之道-从小工到专家》 让浏览器下载文件 再读《软件随想录》/《黑客与画家》/《软技能》 HTTP 压力测试中的 Coordinated Omission 2的补码 编程语言中的 context 是什么? flutter macOS 构建出错 Flatpak 使用小记 Golang CAS 操作是怎么实现的 PostgreSQL 当MQ来使用 Clash 结合 工作VPN 的网络设计 使用 PostgreSQL 搭建 JuiceFS PostgreSQL 配置优化和日志分析 有GitHub Copilot?那就可以搭建你的ChatGPT4服务 窗口函数的使用(以PG为例)
Golang 分布式异步任务队列 Machinery 教程
Jiajun Huang · 2017-12-30 · via Jiajun的技术笔记

Golang的分布式任务队列还不算多,目前比较成熟的应该就只有 Machinery 了。

这篇文章里我们简略的看一下Machinery怎么用。但是我们首先简单介绍一下异步任务这个概念。

如果你熟悉Python中的异步任务框架的话,想必一定听过Celery。异步任务框架是什么呢?异步任务的主要作用是将需要长时间执行 的代码放到一个单独的程序中,例如调用第三方邮件接口,但是这个接口可能非常慢才响应,而你又想确保自己的API及时响应。这个 时候就可以采用异步任务来进行解耦。

一般来说,异步任务都由这么几部分组成:

- broker:broker是用来传递信息的,我们可以想象成“信使”,“外卖配送员”,它的作用是暂时保存产生的任务以便于消费
- 生产者:它负责产生任务
- 消费者:它负责消费任务
- result backend:这个不是必需,但是如果有保存结果的需要,那么就需要它。

而流程则是:

生产者发布任务 -> broker -> 消费者竞争一个任务,然后进行消费 -> (可选:消费后向broker确认已经消费,然后broker删除此任务,
否则将超时重发任务) -> result backend保存结果

Machinery

首先我们来把 Machinery 代码拉下来:

$ go get -u github.com/RichardKnop/machinery/v1

Machinery 对消息的定义是:

// Signature represents a single task invocation
type Signature struct {
	UUID           string
	Name           string
	RoutingKey     string
	ETA            *time.Time
	GroupUUID      string
	GroupTaskCount int
	Args           []Arg
	Headers        Headers
	Immutable      bool
	RetryCount     int
	RetryTimeout   int
	OnSuccess      []*Signature
	OnError        []*Signature
	ChordCallback  *Signature
}

就如同自己写任务队列可能用json一样。

一般生产者先调用 signature := tasks.NewSignature 定义好任务,然后 machineryServer.SendTask 就完成了任务的产生。

Machinery 的异步任务长这样:

func Add(args ...int64) (int64, error) {
  sum := int64(0)
  for _, arg := range args {
    sum += arg
  }
  return sum, nil
}

要注意一点,函数的最后一个参数必需是 error。然后这样注册任务。

server.RegisterTasks(map[string]interface{}{
  "add":      Add,
})

消费者先调用 worker := machineryServer.NewWorker("send_sms", 10) 然后 worker.Launch() 开始监听broker并且消费任务。 当你产生一个任务,名字是 add 时,这个函数就会被调用。

一般你可以把生产者和消费者放到两个文件里,分别定义main函数,然后自己写Makefile,这样就可以直接make然后产生两个可执行 文件,不过我个人更喜欢用 flag 来标识到底是什么身份:

func main() {
	// parse cmd args
	flag.Parse()

	// init config
	initConfig()

	// init machinery worker
	initMachinery()

	// register tasks
	machineryServer.RegisterTask("sendSMS", sendSMS)

	if *worker {
		startWorker()
	} else {
		startWebServer()
	}
}