Sobes.tech
Senior

Documentupdates worden naar de service gestuurd message Document { string Url = 1; // URL van het document, de unieke identificatie uint64 PubDate = 2; // tijd van aangekondigde publicatie van het document uint64 FetchTime = 3; // tijd van ontvangst van deze documentupdate, kan worden beschouwd als versienummer. Het paar (Url, FetchTime) is uniek. string Text = 4; // tekst van het document uint64 FirstFetchTime = 5; // aanvankelijk afwezig, moet worden ingevuld } Documenten kunnen in willekeurige volgorde binnenkomen (niet in de volgorde waarin ze werden bijgewerkt), en er kunnen dubbele berichten zijn. Het is nodig om dezelfde berichten te vormen, maar met gecorrigeerde velden volgens de volgende regels (alles hieronder geldt voor een groep documenten met hetzelfde veld Url): Het veld Text en FetchTime moeten hetzelfde zijn als die van het document met de grootste FetchTime die tot nu toe is ontvangen. Het veld PubDate moet hetzelfde zijn als dat van het bericht met de kleinste FetchTime. Het veld FirstFetchTime moet gelijk zijn aan de minimale waarde van FetchTime. Met andere woorden, op elk moment nemen we PubDate en FirstFetchTime van de eerste versie die tot nu toe is ontvangen (als we ze sorteren op FetchTime), en Text van de laatste. De interface in de code kan als volgt worden geïmplementeerd: type Processor interface { Process(doc *Document) (*Document, error) } Deze code werkt in een service die berichten leest uit een wachtrij (Kafka of vergelijkbaar), en ook het resultaat in de wachtrij schrijft. Als Process null retourneert, wordt er niets in de wachtrij geschreven.

sobes.tech AI

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