Agregando eventos de alto volume a montante (sessioning)

Prev Next

Esta página fornece informações sobre como gerenciar eventos de alto volume a montante e exemplos de implementação para cada plataforma Kafka gerenciada suportada.

Tipos de Agregação

Agregados baseados em sessões

Limitado pelo tempo (por exemplo, ≤10 minutos) ou contagem (por exemplo, ≤N giros). Cada agregado deve incluir totais para apostas, vitórias, net, contagens, código do jogo, canal e carimbos de tempo de início/fim.

Agregados de janelas rolantes

Normalmente, resumos de 2 a 5 minutos agrupados por usuário ou produto.

Por quê EMIT FINAL / DISCARDING A acumulação importa

Agregações de janelas de sessão podem ser ativadas várias vezes à medida que novos eventos chegam dentro de uma sessão aberta. Se seu pipeline usar o comportamento padrão de atualização incremental, sistemas posteriores recebem múltiplos registros parciais para a mesma sessão — fazendo com que a contagem de eventos seja inflada. Os EMIT FINAL padrões (ksqlDB), emissão apenas de janelas fechadas (Fluxos Kafka) e AccumulationMode.DISCARDING (Beam) garantem que exatamente um registro seja emitido por sessão concluída, que é a entrada correta para análises em nível de evento Xtremepush.

Confluent Cloud (ksqlDB)

Agregação de sessões

CREATE TABLE session_summaries
  WITH (KAFKA_TOPIC = 'xp-session-events', VALUE_FORMAT = 'JSON')
AS SELECT
  user_id,
  'session'                                          AS event,
  AS_VALUE(user_id)                                  AS user_id,
  TIMESTAMPTOSTRING(MIN(ROWTIME), 'yyyy-MM-dd HH:mm:ss', 'UTC') AS `timestamp`,
  STRUCT(
    session_start_ms := MIN(ROWTIME),
    session_end_ms   := MAX(ROWTIME),
    event_count      := COUNT(*),
    duration_ms      := MAX(ROWTIME) - MIN(ROWTIME)
  )                                                  AS value,
  STRUCT()                                           AS user_attributes
FROM xtremepush_events
WINDOW SESSION (10 MINUTES)
GROUP BY user_id
EMIT FINAL;

EMIT FINAL é crítico. Sem ele, o ksqlDB emite atualizações incrementais à medida que a sessão está sendo construída. EMIT FINAL garante que exatamente um registro seja emitido por sessão fechada (após o término do intervalo de inatividade), mantendo a contagem de eventos a jusante significativa.

Janela de rodas (resumo de 2 minutos)

CREATE TABLE rolling_2min_summary
  WITH (KAFKA_TOPIC = 'xp-rolling-2min', VALUE_FORMAT = 'JSON')
AS SELECT
  user_id,
  'session_summary'                                  AS event,
  AS_VALUE(user_id)                                  AS user_id,
  TIMESTAMPTOSTRING(WINDOWSTART, 'yyyy-MM-dd HH:mm:ss', 'UTC') AS `timestamp`,
  STRUCT(
    window_start_ms := WINDOWSTART,
    window_end_ms   := WINDOWEND,
    event_count     := COUNT(*)
  )                                                  AS value,
  STRUCT()                                           AS user_attributes
FROM xtremepush_events
WINDOW TUMBLING (SIZE 2 MINUTES)
GROUP BY user_id
EMIT FINAL;

Amazon MSK (Kafka Streams — Java)

Agregação de sessões

import org.apache.kafka.streams.kstream.SessionWindows;
import java.time.Duration;

KTable<Windowed<String>, Long> sessions = rawStream
    .groupByKey()
    .windowedBy(
        SessionWindows.ofInactivityGapAndGrace(
            Duration.ofMinutes(10),   // inactivity gap — closes session after 10min silence
            Duration.ofMinutes(2)     // grace period — accept late-arriving events
        )
    )
    .aggregate(
        () -> 0L,                                      // initialiser
        (key, value, aggregate) -> aggregate + 1,      // incrementor
        (key, agg1, agg2) -> agg1 + agg2              // merger (called when two sessions merge)
    );

sessions.toStream().map((windowed, count) -> {
    try {
        ObjectNode envelope = mapper.createObjectNode();
        envelope.put("event",   "session");
        envelope.put("user_id", windowed.key());

        String ts = Instant.ofEpochMilli(windowed.window().end())
            .toString().replace("T", " ").replaceAll("\\.\\d+Z$", "");
        envelope.put("timestamp", ts);

        ObjectNode value = mapper.createObjectNode();
        value.put("session_start_ms", windowed.window().start());
        value.put("session_end_ms",   windowed.window().end());
        value.put("event_count",      count);
        value.put("duration_ms",      windowed.window().end() - windowed.window().start());
        envelope.set("value", value);
        envelope.set("user_attributes", mapper.createObjectNode());

        return KeyValue.pair(windowed.key(), mapper.writeValueAsString(envelope));
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
}).to("xp-session-events");

A função (key, agg1, agg2) -> agg1 + agg2 de fusão é necessária para janelas de sessão. Os Streams Kafka podem fundir duas janelas de sessão adjacentes quando um evento que chega tardiamente faz a ponte entre elas — a fusão combina as contagens de ambas as janelas.

Janela de rodas (resumo de 2 minutos)

import org.apache.kafka.streams.kstream.TimeWindows;

KTable<Windowed<String>, Long> rolling = rawStream
    .groupByKey()
    .windowedBy(
        TimeWindows.ofSizeAndGrace(Duration.ofMinutes(2), Duration.ofMinutes(1))
    )
    .count();

Google Cloud Managed Apache Kafka (Apache Beam / Dataflow)

Agregação de sessões

import apache_beam as beam
from apache_beam.transforms.window import Sessions
from apache_beam.transforms.trigger import AfterWatermark, AccumulationMode
import json
from datetime import datetime, timezone

class SessionOutputFn(beam.DoFn):
    def process(self, element, window=beam.DoFn.WindowParam):
        user_id, events = element
        ts = datetime.fromtimestamp(window.end, tz=timezone.utc).strftime('%Y-%m-%d %H:%M:%S')
        yield json.dumps({
            'event':   'session',
            'user_id': user_id,
            'timestamp': ts,
            'value': {
                'session_start_ms': int(window.start * 1000),
                'session_end_ms':   int(window.end   * 1000),
                'event_count':      len(list(events)),
                'duration_ms':      int((window.end - window.start) * 1000),
            },
            'user_attributes': {}
        })


(
    events_pcoll
    | 'ExtractUserId'  >> beam.Map(lambda e: (json.loads(e)['user_id'], e))
    | 'SessionWindow'  >> beam.WindowInto(
                              Sessions(gap_size=600),        # 10-minute inactivity gap
                              trigger=AfterWatermark(),
                              accumulation_mode=AccumulationMode.DISCARDING)
    | 'GroupByKey'     >> beam.GroupByKey()
    | 'BuildSession'   >> beam.ParDo(SessionOutputFn())
    | 'WriteToKafka'   >> beam.io.WriteToKafka(
                              producer_config={'bootstrap.servers': 'BROKER:9092'},
                              topic='xp-session-events')
)

AccumulationMode.DISCARDING é crítico. As janelas de viga disparam quando a marca d'água passa pela extremidade da janela. Com DISCARDING, cada painel descarta dados previamente acumulados após disparar — então exatamente um registro completo de sessão é emitido por sessão fechada. Sem esse modo (ACCUMULATING de modo ), você receberia atualizações incrementais de sessão parcial que superariam a contagem dos eventos a jusante.

Janela de rodas (resumo de 2 minutos)

from apache_beam.transforms.window import FixedWindows
from apache_beam.transforms.trigger import AfterWatermark, AccumulationMode

(
    events_pcoll
    | 'KeyByUser'     >> beam.Map(lambda e: (json.loads(e)['user_id'], 1))
    | 'FixedWindow'   >> beam.WindowInto(
                             FixedWindows(120),              # 2-minute fixed windows
                             trigger=AfterWatermark(),
                             accumulation_mode=AccumulationMode.DISCARDING)
    | 'CountPerUser'  >> beam.CombinePerKey(sum)
    | 'BuildSummary'  >> beam.Map(lambda kv: json.dumps({
                             'event':   'session_summary',
                             'user_id': kv[0],
                             'timestamp': datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S'),
                             'value':   {'event_count': kv[1]},
                             'user_attributes': {}
                         }))
    | 'WriteToKafka'  >> beam.io.WriteToKafka(...)
)