Senior
Құжаттар жаңартулары қызметке келеді message Document { string Url = 1; // Құжаттың URL-і, оның бірегей идентификаторы uint64 PubDate = 2; // жариялануы туралы жарияланған уақыт uint64 FetchTime = 3; // осы құжат жаңартуын алу уақыты, оны нұсқау идентификаторы ретінде қарастыруға болады. Жұп (Url, FetchTime) бірегей. string Text = 4; // құжаттың мәтіні uint64 FirstFetchTime = 5; // бастапқыда жоқ, толтыру керек } Құжаттар кез келген тәртіпте келуі мүмкін (олар жаңартылған тәртіпте емес), және қайталанатын хабарламалар болуы мүмкін. Шығарда, бірдей хабарламалар құрастыру керек, бірақ келесі ережелер бойынша түзетілген өрістермен (төмендегі барлық нәрсе бірдей Url өрісі бар құжаттар тобына қатысты): Text және FetchTime өрістері, қазірге дейін алынған ең үлкен FetchTime-ға сәйкес болуы керек. PubDate өрісі, ең кіші FetchTime-ға сәйкес болуы керек. FirstFetchTime өрісі, ең төменгі FetchTime мәніне тең болуы керек. Яғни, әр уақытта FetchTime бойынша сұрыпталған алғашқы нұсқадан PubDate және FirstFetchTime алынады, ал соңғы нұсқадан Text. Кодтағы интерфейс келесідей жүзеге асырылуы мүмкін: type Processor interface { Process(doc *Document) (*Document, error) } Бұл код Kafka немесе ұқсас кезектен хабарламалар оқитын және нәтижесін де кезекке жазатын қызметте жұмыс істейді. Егер Process Null қайтарса, ештеңе кезекке жазылмайды.
sobes.tech AI
Ҷавоб аз AI
Для реализации интерфейса 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 и возвращать скорректированные документы согласно правилам.