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/DISCARDINGA acumulação importaAgregaçõ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 FINALpadrões (ksqlDB), emissão apenas de janelas fechadas (Fluxos Kafka) eAccumulationMode.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(...)
)