Skip to content
 
 

Latest commit

 

History

880 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Asynq logo

简单、可靠且高效的Go分布式任务队列

GoDoc Go Report Card Build Status License: MIT Gitter chat

Asynq是一个Go库,用于队列任务并使用工作者异步处理它们。它由Redis支持,设计为可扩展且易于上手。

Asynq工作原理的高级概述:

  • 客户端将任务放入队列
  • 服务器从队列中提取任务,并为每个任务启动一个工作者goroutine
  • 任务由多个工作者并发处理

任务队列被用作在多台机器之间分配工作的机制。系统可以由多个工作者服务器和代理组成,从而实现高可用性和水平扩展。

用例示例

任务队列图

特性

稳定性和兼容性

状态:该库目前正在大力开发中,API频繁变更。

☝️ 重要提示:当前主版本为零(v0.x.x),以适应快速开发和快速迭代,同时获得用户的早期反馈(欢迎对API提供反馈!)。在v1.0.0发布之前,公共API可能会在没有主版本更新的情况下发生变化。

赞助

如果您在生产环境中使用此包,请考虑赞助该项目以表示您的支持!

快速入门

确保您已安装Go(下载)。支持最新的两个Go版本(参见https://go.dev/dl)。

通过创建一个文件夹并在文件夹内运行go mod init github.com/your/repo(了解更多)来初始化您的项目。然后使用go get命令安装Asynq库:

go get -u github.com/hibiken/asynq

确保您正在本地或从Docker容器运行Redis服务器。需要4.0或更高版本。

接下来,编写一个封装任务创建和任务处理的包。

package tasks

import (
    "context"
    "encoding/json"
    "fmt"
    "log"
    "time"
    "github.com/hibiken/asynq"
)

// 任务类型列表
const (
    TypeEmailDelivery   = "email:deliver"
    TypeImageResize     = "image:resize"
)

type EmailDeliveryPayload struct {
    UserID     int
    TemplateID string
}

type ImageResizePayload struct {
    SourceURL string
}

//----------------------------------------------
// 编写一个NewXXXTask函数来创建任务。
// 任务由类型和有效负载组成。
//----------------------------------------------

func NewEmailDeliveryTask(userID int, tmplID string) (*asynq.Task, error) {
    payload, err := json.Marshal(EmailDeliveryPayload{UserID: userID, TemplateID: tmplID})
    if err != nil {
        return nil, err
    }
    return asynq.NewTask(TypeEmailDelivery, payload), nil
}

func NewImageResizeTask(src string) (*asynq.Task, error) {
    payload, err := json.Marshal(ImageResizePayload{SourceURL: src})
    if err != nil {
        return nil, err
    }
    // 可以将任务选项传递给NewTask,这些选项可以在入队时被覆盖。
    return asynq.NewTask(TypeImageResize, payload, asynq.MaxRetry(5), asynq.Timeout(20 * time.Minute)), nil
}

//---------------------------------------------------------------
// 编写一个HandleXXXTask函数来处理输入任务。
// 注意它满足asynq.HandlerFunc接口。
//
// 处理程序不需要是函数。您可以定义一个满足asynq.Handler接口的类型。
// 请参见下面的示例。
//---------------------------------------------------------------

func HandleEmailDeliveryTask(ctx context.Context, t *asynq.Task) error {
    var p EmailDeliveryPayload
    if err := json.Unmarshal(t.Payload(), &p); err != nil {
        return fmt.Errorf("json.Unmarshal failed: %v: %w", err, asynq.SkipRetry)
    }
    log.Printf("发送电子邮件给用户: user_id=%d, template_id=%s", p.UserID, p.TemplateID)
    // 电子邮件发送代码 ...
    return nil
}

// ImageProcessor实现asynq.Handler接口。
type ImageProcessor struct {
    // ... 结构体的字段
}

func (processor *ImageProcessor) ProcessTask(ctx context.Context, t *asynq.Task) error {
    var p ImageResizePayload
    if err := json.Unmarshal(t.Payload(), &p); err != nil {
        return fmt.Errorf("json.Unmarshal failed: %v: %w", err, asynq.SkipRetry)
    }
    log.Printf("调整图像大小: src=%s", p.SourceURL)
    // 图像调整大小代码 ...
    return nil
}

func NewImageProcessor() *ImageProcessor {
	return &ImageProcessor{}
}

在您的应用程序代码中,导入上述包并使用Client将任务放入队列。

package main

import (
    "log"
    "time"

    "github.com/hibiken/asynq"
    "your/app/package/tasks"
)

const redisAddr = "127.0.0.1:6379"

func main() {
    client := asynq.NewClient(asynq.RedisClientOpt{Addr: redisAddr})
    defer client.Close()

    // ------------------------------------------------------
    // 示例1:将任务入队以立即处理。
    //       使用(*Client).Enqueue方法。
    // ------------------------------------------------------

    task, err := tasks.NewEmailDeliveryTask(42, "some:template:id")
    if err != nil {
        log.Fatalf("无法创建任务: %v", err)
    }
    info, err := client.Enqueue(task)
    if err != nil {
        log.Fatalf("无法将任务入队: %v", err)
    }
    log.Printf("已入队任务: id=%s queue=%s", info.ID, info.Queue)


    // ------------------------------------------------------------
    // 示例2:安排任务在未来处理。
    //       使用ProcessIn或ProcessAt选项。
    // ------------------------------------------------------------

    info, err = client.Enqueue(task, asynq.ProcessIn(24*time.Hour))
    if err != nil {
        log.Fatalf("无法安排任务: %v", err)
    }
    log.Printf("已入队任务: id=%s queue=%s", info.ID, info.Queue)


    // ----------------------------------------------------------------------------
    // 示例3:设置其他选项以调整任务处理行为。
    //       选项包括MaxRetry、Queue、Timeout、Deadline、Unique等。
    // ----------------------------------------------------------------------------

    task, err = tasks.NewImageResizeTask("https://example.com/myassets/image.jpg")
    if err != nil {
        log.Fatalf("无法创建任务: %v", err)
    }
    info, err = client.Enqueue(task, asynq.MaxRetry(10), asynq.Timeout(3 * time.Minute))
    if err != nil {
        log.Fatalf("无法将任务入队: %v", err)
    }
    log.Printf("已入队任务: id=%s queue=%s", info.ID, info.Queue)
}

接下来,启动一个工作者服务器以在后台处理这些任务。要启动后台工作者,请使用Server并提供您的Handler来处理任务。

您可以选择使用ServeMux创建一个处理程序,就像使用net/http Handler一样。

package main

import (
    "log"

    "github.com/hibiken/asynq"
    "your/app/package/tasks"
)

const redisAddr = "127.0.0.1:6379"

func main() {
    srv := asynq.NewServer(
        asynq.RedisClientOpt{Addr: redisAddr},
        asynq.Config{
            // 指定要使用的并发工作者数量
            Concurrency: 10,
            // 可选择指定具有不同优先级的多个队列。
            Queues: map[string]int{
                "critical": 6,
                "default":  3,
                "low":      1,
            },
            // 有关其他配置选项,请参阅godoc
        },
    )

    // mux将类型映射到处理程序
    mux := asynq.NewServeMux()
    mux.HandleFunc(tasks.TypeEmailDelivery, tasks.HandleEmailDeliveryTask)
    mux.Handle(tasks.TypeImageResize, tasks.NewImageProcessor())
    // ...注册其他处理程序...

    if err := srv.Run(mux); err != nil {
        log.Fatalf("无法运行服务器: %v", err)
    }
}

有关该库的更详细介绍,请参阅我们的入门指南

要了解有关asynq功能和API的更多信息,请参阅包godoc

Web UI

Asynqmon是一个基于Web的工具,用于监控和管理Asynq队列和任务。

以下是Web UI的一些截图:

队列视图

Web UI队列视图

任务视图

Web UI任务视图

指标视图 Screen Shot 2021-12-19 at 4 37 19 PM

设置和自适应暗模式

Web UI设置和自适应暗模式

有关如何使用该工具的详细信息,请参阅该工具的README

命令行工具

Asynq附带一个命令行工具,用于检查队列和任务的状态。

要安装CLI工具,请运行以下命令:

go install github.com/hibiken/asynq/tools/asynq@latest

以下是运行asynq dash命令的示例:

Gif

有关如何使用该工具的详细信息,请参阅该工具的README

贡献

我们欢迎并感谢社区做出的任何贡献(GitHub问题/PR、Gitter频道上的反馈等)。

在贡献之前,请参阅贡献指南

许可证

版权所有 (c) 2019-现在 Ken Hibino贡献者Asynq是根据MIT许可证许可的免费开源软件。官方logo由Vic Shóstak创建,并根据Creative Commons许可(CC0 1.0 Universal)分发。

About

简单、可靠、高效的分布式任务队列... 文档/Wiki 已进行汉化处理

Resources

Code of conduct

Contributing

Stars

2 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages