app.conf.task_protocol = 1
from celery_check.celery_app import app
@app.task(queue='gocelery')
def add(a, b):
return a + b
gotasks.add.apply_async(kwargs={'a': 5456, 'b': 2878}, serializer='json', expires=120)
package main
import (
"fmt"
"github.com/gocelery/gocelery"
"time"
"github.com/gomodule/redigo/redis"
)
// exampleAddTask is integer addition task
// with named arguments
type exampleAddTask struct {
a int
b int
}
func (a *exampleAddTask) ParseKwargs(kwargs map[string]interface{}) error {
fmt.Println(233333333)
kwargA, ok := kwargs["a"]
if !ok {
return fmt.Errorf("undefined kwarg a")
}
kwargAFloat, ok := kwargA.(float64)
if !ok {
return fmt.Errorf("malformed kwarg a")
}
a.a = int(kwargAFloat)
kwargB, ok := kwargs["b"]
if !ok {
return fmt.Errorf("undefined kwarg b")
}
kwargBFloat, ok := kwargB.(float64)
if !ok {
return fmt.Errorf("malformed kwarg b")
}
a.b = int(kwargBFloat)
return nil
}
func (a *exampleAddTask) RunTask() (interface{}, error) {
result := a.a + a.b
fmt.Println(result)
return result, nil
}
func Example_workerWithNamedArguments() {
// create redis connection pool
redisPool := &redis.Pool{
Dial: func() (redis.Conn, error) {
c, err := redis.DialURL("redis://:Carizon@1234@10.11.96.81:6379")
if err != nil {
println(234)
return nil, err
}
return c, err
},
}
// initialize celery client
cli, _ := gocelery.NewCeleryClient(
&gocelery.RedisCeleryBroker{
Pool: redisPool,
QueueName: "gocelery",
},
&gocelery.RedisCeleryBackend{Pool: redisPool},
5, // number of workers
)
// register task
cli.Register("celery_check.gotasks.add", &exampleAddTask{})
// start workers (non-blocking call)
cli.StartWorker()
// wait for client request
time.Sleep(100 * time.Second)
// stop workers gracefully (blocking call)
cli.StopWorker()
}
func main() {
Example_workerWithNamedArguments()
}