job

package
v0.0.2 Latest Latest
Warning

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

Go to latest
Published: Nov 11, 2025 License: Apache-2.0 Imports: 18 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func GenerateHostId

func GenerateHostId(length int) string

func GenerateReqID

func GenerateReqID() int64

func NewExecutorCache

func NewExecutorCache() *executorCache

func RunServer

func RunServer(opts *dto.Options, client SnailJobClient, executors map[string]NewJobExecutor, factory LoggerFactory)

Types

type BaseJobExecutor

type BaseJobExecutor struct {
	LocalLogger  *logrus.Entry
	RemoteLogger *logrus.Entry
	// contains filtered or unexported fields
}

func (*BaseJobExecutor) Context

func (executor *BaseJobExecutor) Context() context.Context

func (*BaseJobExecutor) JobExecute

func (executor *BaseJobExecutor) JobExecute(jobContext dto.JobContext)

JobExecute 模板类

type BaseMapJobExecutor

type BaseMapJobExecutor struct {
	BaseJobExecutor
	// contains filtered or unexported fields
}

func (*BaseMapJobExecutor) DoJobExecute

func (executor *BaseMapJobExecutor) DoJobExecute(jobArgs dto.IJobArgs) dto.ExecuteResult

DoJobExecute 模板类

func (*BaseMapJobExecutor) DoMap

func (executor *BaseMapJobExecutor) DoMap(taskList []interface{}, nextTaskName string) (*dto.ExecuteResult, error)

type BaseMapReduceJobExecutor

type BaseMapReduceJobExecutor struct {
	BaseMapJobExecutor
	// contains filtered or unexported fields
}

func (*BaseMapReduceJobExecutor) BindMapReduceExecute

func (executor *BaseMapReduceJobExecutor) BindMapReduceExecute(child MapReduceExecute)

func (*BaseMapReduceJobExecutor) DoJobExecute

func (executor *BaseMapReduceJobExecutor) DoJobExecute(jobArgs dto.IJobArgs) dto.ExecuteResult

DoJobExecute 模板类

func (*BaseMapReduceJobExecutor) DoJobMapExecute

func (executor *BaseMapReduceJobExecutor) DoJobMapExecute(args *dto.MapArgs) dto.ExecuteResult

type DefaultFormatter

type DefaultFormatter struct {
	ForceColors bool
}

func (*DefaultFormatter) Format

func (f *DefaultFormatter) Format(entry *logrus.Entry) ([]byte, error)

type Dispatcher

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

func Init

func Init(client SnailJobClient, executors map[string]NewJobExecutor, factory LoggerFactory) *Dispatcher

func (*Dispatcher) DispatchJob

func (e *Dispatcher) DispatchJob(dispatchJob dto.DispatchJobRequest) dto.Result

func (*Dispatcher) GetExecutor

func (e *Dispatcher) GetExecutor(name string) (IJobExecutor, error)

func (*Dispatcher) Stop

func (e *Dispatcher) Stop(stopJob dto.StopJob) dto.Result

type HookLogService

type HookLogService struct {
	Wg         sync.WaitGroup
	LogEntryCh chan *dto.JobLogTask
	// contains filtered or unexported fields
}

func NewHookLogService

func NewHookLogService(client SnailJobClient) *HookLogService

func (*HookLogService) Init

func (hls *HookLogService) Init()

type IJobExecutor

type IJobExecutor interface {
	JobExecute(context dto.JobContext)
}

IJobExecutor 执行器接口

type JobExecutorFutureCallback

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

type JobStrategy

type JobStrategy interface {
	DoJobExecute(dto.IJobArgs) dto.ExecuteResult
	// contains filtered or unexported methods
}

type LoggerFactory

type LoggerFactory interface {
	GetRemoteLogger(name string, ctx context.Context) *logrus.Entry
	GetLocalLogger(name string) *logrus.Entry
	GetLogRus() *logrus.Logger
	Init(hls *HookLogService)
}

func NewLoggerFactory

func NewLoggerFactory(opts *dto.Options) LoggerFactory

type LoggerHook

type LoggerHook struct {
	Hls *HookLogService
}

func (*LoggerHook) Fire

func (h *LoggerHook) Fire(entry *logrus.Entry) error

func (*LoggerHook) Levels

func (h *LoggerHook) Levels() []logrus.Level

type MapExecute

type MapExecute interface {
	DoJobMapExecute(args *dto.MapArgs) dto.ExecuteResult
	// contains filtered or unexported methods
}

type MapReduceExecute

type MapReduceExecute interface {
	DoReduceExecute(args *dto.ReduceArgs) dto.ExecuteResult
	DoMergeReduceExecute(args *dto.MergeReduceArgs) dto.ExecuteResult
	BindMapReduceExecute(child MapReduceExecute)
}

type NewJobExecutor

type NewJobExecutor func() IJobExecutor

type SafeBuffer

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

线程安全缓冲区

func NewSafeBuffer

func NewSafeBuffer() *SafeBuffer

NewSafeBuffer 返回一个初始化的 SafeBuffer 实例

func (*SafeBuffer) Add

func (sb *SafeBuffer) Add(entry *dto.JobLogTask)

Add 向缓冲区添加一条数据

func (*SafeBuffer) GetAll

func (sb *SafeBuffer) GetAll() []*dto.JobLogTask

GetAll 获取缓冲区中的所有数据并清空缓冲区

func (*SafeBuffer) Len

func (sb *SafeBuffer) Len() int

Len 返回缓冲区中数据的数量

type Server

type Server struct {
	rpc.UnimplementedUnaryRequestServer
	// contains filtered or unexported fields
}

func (*Server) UnaryRequest

func (s *Server) UnaryRequest(_ context.Context, in *rpc.GrpcSnailJobRequest) (*rpc.GrpcResult, error)

UnaryRequest implements snailjob.UnaryRequestServer

type SnailJobClient

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

func NewSnailJobClient

func NewSnailJobClient(opts *dto.Options, factory LoggerFactory) SnailJobClient

func (*SnailJobClient) SendBatchLogReport

func (receiver *SnailJobClient) SendBatchLogReport(payload []*dto.JobLogTask)

func (*SnailJobClient) SendBatchReportMapTask

func (receiver *SnailJobClient) SendBatchReportMapTask(req dto.MapTaskRequest) constant.StatusEnum

func (*SnailJobClient) SendDispatchResult

func (receiver *SnailJobClient) SendDispatchResult(payload interface{})

func (*SnailJobClient) SendHeartbeat

func (receiver *SnailJobClient) SendHeartbeat()

func (*SnailJobClient) SendToServer

func (receiver *SnailJobClient) SendToServer(uri string, payload interface{}, jobName string) constant.StatusEnum

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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