async

package
v0.0.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 2, 2026 License: MIT Imports: 4 Imported by: 0

README

async 包 — 异步方法执行

所属层级: Infrastructure Layer
设计理念: 简化异步编程,Future 模式
设计灵感: Spring @Async + CompletableFuture

概述

async 包提供异步方法执行功能,参考 Spring Boot 的 @Async 注解设计。基于 goroutine 池实现异步任务执行,支持 Future 模式返回值、优雅关闭等特性。

核心功能
功能 说明
线程池配置 支持核心线程数、最大线程数、队列容量配置
Future 返回值 支持阻塞获取结果、超时获取、状态检查
自定义线程名 支持自定义线程名前缀
优雅关闭 支持带超时的优雅关闭
零依赖 仅使用 Go 标准库

核心接口

Future 异步结果
type Future struct {
    // ...
}
获取结果
方法 说明
Get() 阻塞获取结果
GetWithContext(ctx) 带上下文的阻塞获取
GetWithTimeout(timeout) 带超时的阻塞获取
IsDone() 检查任务是否完成
// 阻塞获取
result, err := future.Get()

// 带超时获取
result, err := future.GetWithTimeout(5 * time.Second)

// 带上下文获取
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
result, err := future.GetWithContext(ctx)

// 状态检查
if future.IsDone() {
    result, err := future.Get()
}
AsyncExecutor 异步执行器
type AsyncExecutor struct {
    // ...
}
创建
executor := async.NewAsyncExecutor(
    4,   // 工作线程数
    100, // 任务队列容量
)
生命周期管理
方法 说明
Start() 启动执行器
Shutdown() 优雅关闭执行器
ShutdownWithTimeout(timeout) 带超时的优雅关闭
IsRunning() 检查执行器是否运行
// 启动
executor.Start()

// 优雅关闭
executor.Shutdown()

// 带超时关闭
err := executor.ShutdownWithTimeout(10 * time.Second)
提交任务
方法 说明
Submit(fn) 提交有返回值的异步任务
SubmitVoid(fn) 提交无返回值的异步任务
GetQueueSize() 获取当前队列中的任务数
// 提交有返回值的任务
future := executor.Submit(func() (any, error) {
    result := doSomething()
    return result, nil
})

// 提交无返回值的任务
executor.SubmitVoid(func() error {
    return doSomething()
})

快速开始

基本使用
package main

import (
    "fmt"
    "github.com/xudefa/enhance/async"
)

func main() {
    executor := async.NewAsyncExecutor(4, 100)
    defer executor.Shutdown()

    // 提交异步任务
    future := executor.Submit(func() (any, error) {
        // 模拟耗时操作
        time.Sleep(1 * time.Second)
        return "result", nil
    })

    // 阻塞获取结果
    result, err := future.Get()
    if err != nil {
        fmt.Println("Error:", err)
        return
    }

    fmt.Println("Result:", result)
}

API 参考

带超时获取
executor := async.NewAsyncExecutor(4, 100)
defer executor.Shutdown()

future := executor.Submit(func() (any, error) {
    time.Sleep(5 * time.Second)
    return "slow result", nil
})

// 设置 3 秒超时
result, err := future.GetWithTimeout(3 * time.Second)
if err != nil {
    fmt.Println("Timeout:", err)
    return
}
使用 Context
executor := async.NewAsyncExecutor(4, 100)
defer executor.Shutdown()

future := executor.Submit(func() (any, error) {
    return fetchDataFromAPI()
})

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()

result, err := future.GetWithContext(ctx)
if err != nil {
    if errors.Is(err, context.DeadlineExceeded) {
        fmt.Println("Request timeout")
    }
    return
}
批量提交任务
executor := async.NewAsyncExecutor(4, 100)
defer executor.Shutdown()

// 提交多个任务
var futures []*async.Future
for i := 0; i < 10; i++ {
    taskID := i
    future := executor.Submit(func() (any, error) {
        return processTask(taskID)
    })
    futures = append(futures, future)
}

// 收集所有结果
for i, future := range futures {
    result, err := future.Get()
    if err != nil {
        fmt.Printf("Task %d failed: %v\n", i, err)
        continue
    }
    fmt.Printf("Task %d result: %v\n", i, result)
}
无返回值任务
executor := async.NewAsyncExecutor(4, 100)
defer executor.Shutdown()

// 提交后台任务
executor.SubmitVoid(func() error {
    return sendEmailNotification("user@example.com", "Welcome!")
})

executor.SubmitVoid(func() error {
    return updateCache("user:123", userData)
})

使用示例

场景 1: 异步邮件发送

用户注册后异步发送邮件通知,不阻塞主流程:

func (s *UserService) RegisterUser(req *RegisterRequest) error {
    // 创建用户
    user := s.createUser(req)

    // 异步发送邮件
    s.executor.SubmitVoid(func() error {
        return s.emailService.SendWelcomeEmail(user.Email)
    })

    // 立即返回
    return nil
}
场景 2: 批量数据处理

批量处理数据时,使用异步提高处理效率:

func (s *DataService) ProcessBatch(items []Item) ([]Result, error) {
    var futures []*async.Future

    // 提交所有处理任务
    for _, item := range items {
        item := item // 捕获循环变量
        future := s.executor.Submit(func() (any, error) {
            return s.processItem(item)
        })
        futures = append(futures, future)
    }

    // 收集结果
    results := make([]Result, 0, len(items))
    for _, future := range futures {
        result, err := future.Get()
        if err != nil {
            return nil, err
        }
        results = append(results, result.(Result))
    }

    return results, nil
}
场景 3: 后台定时任务

执行后台定时任务,如数据清理、缓存刷新等:

func (s *CleanupService) StartCleanupJob() {
    ticker := time.NewTicker(1 * time.Hour)
    defer ticker.Stop()

    for range ticker.C {
        s.executor.SubmitVoid(func() error {
            count, err := s.cleanupExpiredRecords()
            if err != nil {
                log.Printf("Cleanup failed: %v", err)
                return err
            }
            log.Printf("Cleaned up %d records", count)
            return nil
        })
    }
}
场景 4: 并发 API 调用

并发调用多个外部 API,减少总耗时:

func (s *AggregationService) GetAggregatedData(userID string) (*AggregatedData, error) {
    // 并发调用多个 API
    userFuture := s.executor.Submit(func() (any, error) {
        return s.userAPI.GetUser(userID)
    })

    ordersFuture := s.executor.Submit(func() (any, error) {
        return s.orderAPI.GetOrders(userID)
    })

    statsFuture := s.executor.Submit(func() (any, error) {
        return s.statsAPI.GetUserStats(userID)
    })

    // 等待所有结果
    user, err := userFuture.Get()
    if err != nil {
        return nil, err
    }

    orders, err := ordersFuture.Get()
    if err != nil {
        return nil, err
    }

    stats, err := statsFuture.Get()
    if err != nil {
        return nil, err
    }

    return &AggregatedData{
        User:   user.(*User),
        Orders: orders.([]*Order),
        Stats:  stats.(*UserStats),
    }, nil
}

最佳实践

1. 非关键路径操作使用异步
// ✅ 推荐:邮件发送使用异步
s.executor.SubmitVoid(func() error {
    return s.emailService.SendWelcomeEmail(user.Email)
})

// ⚠️ 不推荐:阻塞主流程
err := s.emailService.SendWelcomeEmail(user.Email)
2. 捕获循环变量避免闭包问题
// ✅ 推荐:捕获循环变量
for _, item := range items {
    item := item // 捕获循环变量
    executor.Submit(func() (any, error) {
        return processItem(item)
    })
}

// ⚠️ 不推荐:直接使用循环变量
for _, item := range items {
    executor.Submit(func() (any, error) {
        return processItem(item) // 可能获取到错误的值
    })
}
3. 设置合理的超时时间
// ✅ 推荐:设置超时
result, err := future.GetWithTimeout(5 * time.Second)

// ⚠️ 不推荐:无限期等待
result, err := future.Get()
4. 优雅关闭执行器
// ✅ 推荐:使用 defer 确保关闭
executor := async.NewAsyncExecutor(4, 100)
defer executor.ShutdownWithTimeout(10 * time.Second)

// ⚠️ 不推荐:忘记关闭
executor := async.NewAsyncExecutor(4, 100)
5. 与依赖注入集成
// ✅ 推荐:将 Executor 注册为 Bean
container.Register(
    reflect.TypeOf(&async.AsyncExecutor{}),
    core.Bean(async.NewAsyncExecutor(4, 100)),
    core.Singleton(),
)

// 注入使用
type UserService struct {
    Executor *async.AsyncExecutor `inject:"asyncExecutor"`
}

设计要点

  • AsyncExecutor 使用 context.Context 控制生命周期
  • Future 使用 chan struct{} 实现阻塞等待
  • 任务队列使用 chan 实现,支持缓冲
  • 优雅关闭时等待所有任务完成
  • 零外部依赖,仅使用 Go 标准库

Documentation

Overview

Package async 提供异步执行器功能,用于 enhance 框架。

该模块提供异步任务执行、线程池管理、异步结果获取等异步编程支持。 适用于需要并发执行任务的场景。

架构设计

  • AsyncExecutor: 异步任务执行器,基于 goroutine 池实现
  • Future: 异步结果,支持阻塞获取、超时获取、上下文控制
  • ExecutorOption: 执行器配置选项函数
  • RejectHandler: 任务拒绝处理函数

核心功能

  • 异步任务执行: 基于 goroutine 池的异步任务执行
  • Future 模式: 支持异步结果获取
  • 超时控制: 支持带超时的结果获取
  • 上下文控制: 支持通过上下文取消任务
  • 优雅关闭: 支持执行器的优雅关闭

使用方式

// 创建执行器
executor := async.NewAsyncExecutor(context.Background(), 10, 100)
executor.Start()

// 提交异步任务
future := executor.Submit(func() (any, error) {
    // 执行耗时操作
    return result, nil
})

// 获取结果
result, err := future.Get()

// 带超时获取
result, err := future.GetWithTimeout(5 * time.Second)

// 优雅关闭
executor.Shutdown()

配置选项

  • WithPoolSize: 设置线程池大小
  • WithQueueSize: 设置任务队列大小
  • WithRejectHandler: 设置任务拒绝策略

Package async 提供异步执行器功能,用于 enhance 框架。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type AsyncExecutor

type AsyncExecutor struct {
	// contains filtered or unexported fields
}

AsyncExecutor 异步执行器。

基于 goroutine 池实现异步任务执行,支持 Future 模式返回值。

func NewAsyncExecutor

func NewAsyncExecutor(ctx context.Context, workerCount, queueSize int) *AsyncExecutor

NewAsyncExecutor 创建异步执行器

参数:

  • ctx: 父级 context,用于控制执行器生命周期
  • workerCount: 工作线程数
  • queueSize: 任务队列容量

func (*AsyncExecutor) GetQueueSize

func (e *AsyncExecutor) GetQueueSize() int

GetQueueSize 获取当前队列中的任务数

func (*AsyncExecutor) IsRunning

func (e *AsyncExecutor) IsRunning() bool

IsRunning 检查执行器是否运行

func (*AsyncExecutor) Shutdown

func (e *AsyncExecutor) Shutdown()

Shutdown 优雅关闭执行器

func (*AsyncExecutor) ShutdownWithTimeout

func (e *AsyncExecutor) ShutdownWithTimeout(timeout time.Duration) error

ShutdownWithTimeout 带超时的优雅关闭

func (*AsyncExecutor) Start

func (e *AsyncExecutor) Start()

Start 启动执行器

func (*AsyncExecutor) Submit

func (e *AsyncExecutor) Submit(fn func() (any, error)) *Future

Submit 提交异步任务

func (*AsyncExecutor) SubmitVoid

func (e *AsyncExecutor) SubmitVoid(fn func() error) *Future

SubmitVoid 提交无返回值的异步任务

type ExecutorOption

type ExecutorOption func(*AsyncExecutor)

ExecutorOption 执行器配置选项函数。

type Future

type Future struct {
	// contains filtered or unexported fields
}

Future 异步任务结果。

提供阻塞获取结果、超时获取、检查是否完成等功能。

func NewFuture

func NewFuture() *Future

NewFuture 创建 Future 实例

func (*Future) Get

func (f *Future) Get() (any, error)

Get 阻塞获取结果

func (*Future) GetWithContext

func (f *Future) GetWithContext(ctx context.Context) (any, error)

GetWithContext 带上下文的阻塞获取

func (*Future) GetWithTimeout

func (f *Future) GetWithTimeout(timeout time.Duration) (any, error)

GetWithTimeout 带超时的阻塞获取

func (*Future) IsDone

func (f *Future) IsDone() bool

IsDone 检查任务是否完成

type RejectHandler

type RejectHandler func(task func() (any, error))

RejectHandler 任务拒绝处理函数类型。

当任务队列满时调用,用于处理任务被拒绝的情况。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL