Senior
Hujjatlar bo'yicha yangilanishlar xizmatga keladi message Document { string Url = 1; // Hujjat URL, uning noyob identifikatori uint64 PubDate = 2; // e'lon qilingan vaqti uint64 FetchTime = 3; // ushbu hujjat yangilanishining qabul qilish vaqti, versiya identifikatori sifatida qaralishi mumkin. Juftlik (Url, FetchTime) noyob. string Text = 4; // hujjat matni uint64 FirstFetchTime = 5; // dastlab mavjud emas, to'ldirilishi kerak } Hujjatlar istalgan tartibda kelishi mumkin (yangilashlar qachon amalga oshirilganiga qarab emas), shuningdek, takroriy xabarlar bo'lishi mumkin. Chiqishda, shunday xabarlar shakllantirish kerakki, ular bilan bir xil bo'lsa, lekin quyidagi qoidalarga muvofiq tuzatilgan maydonlar bilan (quyida aytilganlar, Url maydoniga mos keladigan hujjatlar guruhiga taalluqlidir): Text va FetchTime maydonlari, hozirgacha olingan eng katta FetchTime bilan bir xil bo'lishi kerak. PubDate maydoni, eng kichik FetchTime bilan bir xil bo'lishi kerak. FirstFetchTime maydoni, eng kichik FetchTime qiymatiga teng bo'lishi kerak. Boshqacha qilib aytganda, har doim FetchTime bo'yicha saralangan birinchi versiyadan PubDate va FirstFetchTime, va oxirgi versiyadan Text olinadi. Kodni quyidagicha amalga oshirish mumkin: type Processor interface { Process(doc *Document) (*Document, error) } Ushbu kod, Kafka yoki shunga o'xshash navbatdan xabarlarni o'qiydigan va natijani ham navbatga yozadigan xizmatda ishlaydi. Agar Process Null qaytarilsa, hech narsa navbatga yozilmaydi.
sobes.tech AI
AIdan javob
Для реализации интерфейса Processor в Go нужно хранить состояние по каждому Url, чтобы на каждый входящий документ обновлять агрегированное состояние и возвращать скорректированное сообщение.
Основная идея:
- Для каждого Url хранить:
- минимальный FetchTime (для PubDate и FirstFetchTime)
- максимальный FetchTime (для Text и FetchTime)
- PubDate, Text, FirstFetchTime, соответствующие этим минимальному и максимальному FetchTime
При получении нового документа:
- Если для Url нет записи — создать с данными из документа
- Иначе обновить минимальный и максимальный FetchTime и соответствующие поля при необходимости
- Вернуть агрегированный документ с нужными полями
Пример реализации:
import (
"errors"
"sync"
)
type Document struct {
Url string
PubDate uint64
FetchTime uint64
Text string
FirstFetchTime uint64
}
type Processor interface {
Process(doc *Document) (*Document, error)
}
type docState struct {
minFetchTime uint64
maxFetchTime uint64
pubDate uint64
text string
firstFetch uint64
}
type processorImpl struct {
mu sync.Mutex
state map[string]*docState
}
func NewProcessor() Processor {
return &processorImpl{
state: make(map[string]*docState),
}
}
func (p *processorImpl) Process(doc *Document) (*Document, error) {
if doc == nil {
return nil, errors.New("input document is nil")
}
p.mu.Lock()
defer p.mu.Unlock()
st, ok := p.state[doc.Url]
if !ok {
// первый документ для данного Url
p.state[doc.Url] = &docState{
minFetchTime: doc.FetchTime,
maxFetchTime: doc.FetchTime,
pubDate: doc.PubDate,
text: doc.Text,
firstFetch: doc.FetchTime,
}
return &Document{
Url: doc.Url,
PubDate: doc.PubDate,
FetchTime: doc.FetchTime,
Text: doc.Text,
FirstFetchTime: doc.FetchTime,
}, nil
}
updated := false
// Обновляем минимальный FetchTime и PubDate, FirstFetchTime
if doc.FetchTime < st.minFetchTime {
st.minFetchTime = doc.FetchTime
st.pubDate = doc.PubDate
st.firstFetch = doc.FetchTime
updated = true
}
// Обновляем максимальный FetchTime и Text
if doc.FetchTime > st.maxFetchTime {
st.maxFetchTime = doc.FetchTime
st.text = doc.Text
updated = true
}
// Если пришло сообщение с FetchTime, который уже есть, но текст или PubDate отличаются,
// можно игнорировать, так как пара (Url, FetchTime) уникальна, либо обновлять по логике.
if !updated {
// Данные не изменились, можно вернуть nil, чтобы не писать в очередь
return nil, nil
}
return &Document{
Url: doc.Url,
PubDate: st.pubDate,
FetchTime: st.maxFetchTime,
Text: st.text,
FirstFetchTime: st.firstFetch,
}, nil
}
Таким образом, сервис будет поддерживать актуальное агрегированное состояние для каждого Url и возвращать скорректированные документы согласно правилам.