perf(redisstore): batch subscriber settings into single HGETALL

ListSubscribers fetched settings with one HGET per subscriber. Replace
the per-key fetch with a single HGETALL and resolve each subscriber from
an in-memory map, cutting N round-trips to one per delivery cycle.

Extract resolveSubscriberSettings so the batched and single-key paths
share identical decode/default/self-heal behavior.
This commit is contained in:
tiennm99 committed 2026-06-26 09:04:12 +07:00
1 parent 0404e0dd2d
commit 08b99ad68b
2 files changed
+31 -5

No files matched your search

+19 -4
View File
@@ -97,15 +97,30 @@ func (s *Store) ListSubscribers(ctx context.Context) ([]Subscriber, error) {
return nil, err return nil, err
} }
// Fetch every subscriber's settings in one HGETALL instead of one HGET per
// key. The settings hash mirrors the subscribers set, so its size is the
// subscriber count rather than unbounded history.
settingsByKey, err := s.client.HGetAll(ctx, subscriberSettingsKey).Result()
if err != nil {
return nil, err
}
subscribers := make([]Subscriber, 0, len(keys)) subscribers := make([]Subscriber, 0, len(keys))
for _, key := range keys { for _, key := range keys {
sub, ok, err := s.loadSubscriber(ctx, key) sub, err := ParseSubscriberKey(key)
if err != nil {
if removeErr := s.client.SRem(ctx, subscribersKey, key).Err(); removeErr != nil {
return nil, fmt.Errorf("remove malformed subscriber key %q: %w", key, removeErr)
}
return nil, fmt.Errorf("malformed subscriber key %q: %w", key, err)
}
value, present := settingsByKey[key]
settings, err := s.resolveSubscriberSettings(ctx, key, value, present)
if err != nil { if err != nil {
return nil, err return nil, err
} }
if !ok { sub.Types = settings.Types
continue sub.Components = settings.Components
}
subscribers = append(subscribers, sub) subscribers = append(subscribers, sub)
} }
return subscribers, nil return subscribers, nil
+12 -1
View File
@@ -33,12 +33,23 @@ func (s *Store) loadExistingSubscriberSettings(ctx context.Context, key string)
func (s *Store) loadSubscriberSettings(ctx context.Context, key string) (subscriberSettings, error) { func (s *Store) loadSubscriberSettings(ctx context.Context, key string) (subscriberSettings, error) {
value, err := s.client.HGet(ctx, subscriberSettingsKey, key).Result() value, err := s.client.HGet(ctx, subscriberSettingsKey, key).Result()
if err == redis.Nil { if err == redis.Nil {
return defaultSubscriberSettings(), nil return s.resolveSubscriberSettings(ctx, key, "", false)
} }
if err != nil { if err != nil {
return subscriberSettings{}, err return subscriberSettings{}, err
} }
return s.resolveSubscriberSettings(ctx, key, value, true)
}
// resolveSubscriberSettings turns a raw settings hash value into decoded
// settings. It lets callers that already hold the value (e.g. a batched
// HGETALL in ListSubscribers) reuse the same decode/default/self-heal logic as
// the single-key HGET path. Absent or corrupt values fall back to defaults,
// dropping the corrupt field so it heals on next write.
func (s *Store) resolveSubscriberSettings(ctx context.Context, key, value string, present bool) (subscriberSettings, error) {
if !present {
return defaultSubscriberSettings(), nil
}
settings, err := decodeSubscriberSettings(key, value) settings, err := decodeSubscriberSettings(key, value)
if err == nil { if err == nil {
return settings, nil return settings, nil