Sobes.tech
Senior

Dokumentų atnaujinimai atkeliauja į paslaugą message Document { string Url = 1; // Dokumento URL, jo unikalus identifikatorius uint64 PubDate = 2; // deklaruotas dokumento paskelbimo laikas uint64 FetchTime = 3; // šio dokumento atnaujinimo gavimo laikas, gali būti laikomas versijos identifikatoriumi. Poros (Url, FetchTime) unikalus. string Text = 4; // dokumento tekstas uint64 FirstFetchTime = 5; // pradinė reikšmė nėra, turi būti užpildyta } Dokumentai gali būti gauti bet kokia tvarka (ne ta, kuria jie buvo atnaujinti), ir gali būti dubliuotų žinučių. Išėjime reikia sukurti tas pačias žinutes, bet su pataisytais laukais pagal šias taisykles (visa žemiau nurodyta taikoma dokumentų grupei su tuo pačiu Url lauku): Text ir FetchTime laukų reikšmės turi būti tokios kaip didžiausio iki šiol gauto FetchTime dokumento. PubDate turi būti tokia kaip mažiausio FetchTime žinutės. FirstFetchTime turi būti lygi mažiausiai FetchTime reikšmei. Kiekvienu momentu mes paimame PubDate ir FirstFetchTime iš pirmos iki šiol gautos versijos (jei jas rūšiuotume pagal FetchTime), o Text – iš paskutinės. Kodo interfeisas gali būti įgyvendintas taip: type Processor interface { Process(doc *Document) (*Document, error) } Šis kodas veiks paslaugoje, kuri skaito žinutes iš eilės (Kafka ar panašiai), ir taip pat rašo rezultatą į eilę. Jei Process grąžina Null, nieko į eilę neįrašoma.

sobes.tech AI

Atsakymas iš 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 и возвращать скорректированные документы согласно правилам.