Sobes.tech
Senior

Sissepääsuteenusele saabuvad dokumentide uuendused message Document { string Url = 1; // Dokumendi URL, selle unikaalne identifikaator uint64 PubDate = 2; // dokumendi avaldamise aeg uint64 FetchTime = 3; // selle uuenduse saamise aeg, võib käsitleda versiooni identifikaatorina. Paari (Url, FetchTime) unikaalne. string Text = 4; // dokumendi tekst uint64 FirstFetchTime = 5; // algselt puudub, tuleb täita } Dokumente võib saabuda mis tahes järjekorras (mitte nende uuendamise järjekorras), samuti võib esineda duplikaate eraldi sõnumitest. Väljundil tuleb koostada samasugused sõnumid, kuid parandatud väljadega vastavalt järgmistele reeglitele (kõik allpool olevad reeglid kehtivad dokumentide rühma kohta, millel on sama väli Url): Väli Text ja FetchTime peavad olema sellised, nagu olid dokumentides, mille FetchTime oli suurim selle hetke jooksul. Väli PubDate peab olema selline, nagu oli sõnumil, mille FetchTime oli väikseim. Väli FirstFetchTime peab olema võrdne väikseima FetchTime väärtusega. Ehk teisisõnu, iga hetk võtame PubDate ja FirstFetchTime kõige esimese (sorteeritud FetchTime järgi) saadud versiooni järgi, ning Text on viimasest. Koodis saab selle rakendada järgmiselt: type Processor interface { Process(doc *Document) (*Document, error) } See kood töötab teenuses, mis loeb sisendsõnumeid sõnumijärjestusest (näiteks Kafka) ning kirjutab tulemuse samuti järjesele. Kui Process tagastab Null, siis ei kirjutata midagi järjesele.

sobes.tech AI

Vastus AI-lt

Для реализации интерфейса 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 и возвращать скорректированные документы согласно правилам.