mirror of
https://github.com/christianselig/apollo-backend
synced 2024-11-22 03:37:43 +00:00
enqueue directly from redis
This commit is contained in:
parent
4a92b90df9
commit
c4a0704883
1 changed files with 14 additions and 11 deletions
|
@ -144,7 +144,7 @@ func enqueueAccounts(ctx context.Context, logger *logrus.Logger, statsd *statsd.
|
|||
last_enqueued_at < $1
|
||||
OR last_checked_at < $2
|
||||
ORDER BY last_checked_at
|
||||
LIMIT 500
|
||||
LIMIT 400
|
||||
)
|
||||
UPDATE accounts
|
||||
SET last_enqueued_at = $3
|
||||
|
@ -198,6 +198,7 @@ func enqueueAccounts(ctx context.Context, logger *logrus.Logger, statsd *statsd.
|
|||
for i=1, #ids do
|
||||
local key = "locks:accounts:" .. ids[i]
|
||||
if redis.call("exists", key) == 0 then
|
||||
redis.call("lpush", "rmq::queue::[notifications]::ready", ids[i])
|
||||
redis.call("setex", key, %d, 1)
|
||||
retv[#retv + 1] = ids[i]
|
||||
end
|
||||
|
@ -218,17 +219,19 @@ func enqueueAccounts(ctx context.Context, logger *logrus.Logger, statsd *statsd.
|
|||
enqueued += len(vals)
|
||||
skipped += len(batch) - len(vals)
|
||||
|
||||
for _, val := range vals {
|
||||
id := val.(int64)
|
||||
payload := fmt.Sprintf("%d", id)
|
||||
if err = queue.Publish(payload); err != nil {
|
||||
logger.WithFields(logrus.Fields{
|
||||
"accountID": payload,
|
||||
"err": err,
|
||||
}).Error("failed to enqueue account")
|
||||
failed++
|
||||
/*
|
||||
for _, val := range vals {
|
||||
id := val.(int64)
|
||||
payload := fmt.Sprintf("%d", id)
|
||||
if err = queue.Publish(payload); err != nil {
|
||||
logger.WithFields(logrus.Fields{
|
||||
"accountID": payload,
|
||||
"err": err,
|
||||
}).Error("failed to enqueue account")
|
||||
failed++
|
||||
}
|
||||
}
|
||||
}
|
||||
*/
|
||||
}
|
||||
|
||||
statsd.Histogram("apollo.queue.enqueued", float64(enqueued), []string{}, 1)
|
||||
|
|
Loading…
Reference in a new issue