-
Notifications
You must be signed in to change notification settings - Fork 31
/
redis_stream_reader.go
52 lines (46 loc) · 1.39 KB
/
redis_stream_reader.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
// Copyright (c) 2020 Zededa, Inc.
// SPDX-License-Identifier: Apache-2.0
package redis
import (
"bytes"
"errors"
"io"
"time"
"github.com/go-redis/redis"
)
// RedisStreamReader reads msgpack messages from Redis streams and turns then into JSON strings
type RedisStreamReader struct {
// Redis client handle
Client *redis.Client
// Name of a stream
Stream string
// offset
offset string
}
// Next returns reader for the next chunk of data (message), its size and possible error
func (d *RedisStreamReader) Next() (io.Reader, int64, error) {
if d.Client == nil || d.Stream == "" {
return nil, 0, errors.New("redis connection and name of the stream required")
}
if d.offset == "" {
d.offset = "0"
}
records, err := d.Client.XRead(&redis.XReadArgs{
Streams: []string{d.Stream, d.offset},
Block: time.Millisecond, // do a non-blocking read
Count: 1,
}).Result()
// it is weird that the library would return "redis: nil" for a non-blocking read
if (err != nil && err.Error() != "redis: nil") || len(records) > 1 {
return nil, 0, errors.New("failed to read from stream")
}
if records == nil || len(records[0].Messages) == 0 {
return nil, 0, nil
}
d.offset = records[0].Messages[0].ID
s, ok := records[0].Messages[0].Values["object"].(string)
if !ok {
return nil, 0, errors.New("failed to read from stream")
}
return bytes.NewReader([]byte(s)), int64(len([]byte(s))), nil
}