Formatos de tópicos e requisitos de carga útil

Prev Next

Informações sobre tópicos, formatos e cargas úteis disponíveis sobre kafka, com exemplos.

Todas as cargas úteis devem ser JSON válidas. Cada projeto tem dois temas, cada um com um propósito distinto.

Tema dos eventos

Use este tópico para registrar eventos e ações que um usuário realiza (por exemplo, apostas, depósitos, logins).

  • Nome do tópico: Definido como parte da configuração

  • Chave da mensagem: Não é obrigatório, mas pode ser definido para um UUID único de evento. O Xtremepush irá deduplicar eventos dentro de uma janela de tempo configurável com base nessa chave. Para permitir a deduplicação do seu projeto, entre em contato com o suporte do Xtremepush.

  • Cabeçalhos das mensagens: Nenhum

  • Corpo da mensagem: Item codificado em JSON

Propriedade

Tipo

Obrigatório

Descrição

event

Corda

Sim

O nome do evento, por exemplo, bet. deposit

user_id

Corda

Pelo menos um identificador de usuário é necessário

ID de usuário único. Um novo usuário será criado automaticamente se não houver nenhum.

customer_id

Corda

Pelo menos um identificador de usuário exige

Identificador adicional de usuário, se disponível.

profile_id

Corda

Pelo menos um identificador de usuário exige

Identificador de perfil Xtremepush.

device_id

Corda

Pelo menos um identificador de usuário exige

Identificador de dispositivo Xtremepush.

user_attributes

Objetivo

Não

Informações adicionais sobre usuários. Os atributos são adicionados automaticamente ao perfil do usuário durante o processamento de eventos.

value

Objetivo

Não

Propriedades de evento específicas para o tipo de evento. Pode conter qualquer número de pares-chave-valor, incluindo arrays aninhados ou objetos.

timestamp

Corda

Sim

O carimbo de data original do evento.

Esquema

{
  "event": "some_event",
  "user_id": "some_user",
  "user_attributes": {
    "any_attr": "any_value"
  },
  "value": {
    "any_key": "any_value"
  },
  "timestamp": "2025-01-01 12:00:00"
}

Exemplo

{
  "event": "bet",
  "user_id": "3e685367-07d5-4d48-93ae-f007ac336605",
  "customer_id": 12345,
  "user_attributes": {
    "customer_tier": "VIP"
  },
  "value": {
    "bet_id": "3e685367-07d5-4d48-93ae-f007ac336605",
    "odds": 12.2,
    "stake": 100.0
  },
  "timestamp": "2024-09-01 12:00:00.123"
}

Tema dos usuários

Use este tópico para criar novos usuários ou atualizar perfis de usuários existentes.

Importante — sequenciamento

Mensagens de perfil são processadas sequencialmente e importações grandes podem levar um tempo significativo para serem concluídas.

  • Nome do tópico: Definido como parte da configuração

  • Chave da mensagem: Recomendado. Deve ser o identificador de usuário (por exemplo, user_id). Usado para manter a ordem dentro do fluxo Kafka — todas as mensagens com a mesma chave são roteadas para a mesma partição, garantindo que as atualizações de perfil sejam aplicadas na sequência correta.

  • Cabeçalhos das mensagens: Nenhum

  • Corpo da mensagem: Item codificado em JSON

Propriedade

Tipo

Obrigatório

Descrição

user_id

Corda

Sim

ID de usuário único. Um novo usuário será criado automaticamente se não houver nenhum.

user_attributes

Objetivo

Sim

Informações sobre o usuário. Os atributos são adicionados automaticamente ao perfil do usuário durante o processamento.

customer_id

Corda

Não

Identificador adicional de usuário, se disponível.

Carimbo de tempo

Corda

Não

O carimbo de data e hora da mudança de atributo. Usado para garantir que o valor mais recente seja economizado. Ele será usado por padrão no carimbo de data e hora de uma mensagem no Kafka se estiver ausente.

Esquema

{
  "user_id": "some_user",
  "user_attributes": {
    "any_attr": "any_value"
  },
  "timestamp": "2025-01-01 12:00:00"
}

Exemplo

{
  "user_id": "3e685367-07d5-4d48-93ae-f007ac336605",
  "customer_id": 12345,
  "user_attributes": {
    "customer_tier": "VIP",
    "email": "[email protected]"
  },
  "timestamp": "2024-09-01 12:00:00.123"
}

Transformando seus eventos para o formato Xtremepush

Se seus eventos forem produzidos em um esquema plano e específico para o cliente, você precisará transformá-los no formato padrão de envelope Xtremepush antes (ou como) eles forem publicados no tema dos eventos.

Princípio chave: A transformação é apenas estrutural. Você envolve suas propriedades de eventos existentes no objeto padrão value: {} e adiciona os campos de envelope externo (event, user_id, timestamp, user_attributes). Seus nomes de chaves de propriedade interna (por exemplo, AmountWagered, BetId) são preservados exatamente como estão — você não precisa renomeá-los.

Antes e depois

Antes (mensagem crua e plana do cliente sobre Kafka):

{
  "event_type":    "bet_placed",
  "user_id":       "usr_123",
  "timestamp":     "2025-01-15T14:23:45.000Z",
  "AmountWagered": 50.00,
  "BetId":         "bet_abc_789",
  "SportName":     "Football"
}

Após (envelope padrão Xtremepush):

{
  "event":     "bet_placed",
  "user_id":   "usr_123",
  "timestamp": "2025-01-15 14:23:45",
  "value": {
    "AmountWagered": 50.00,
    "BetId":         "bet_abc_789",
    "SportName":     "Football"
  },
  "user_attributes": {}
}

AmountWagered, BetId, e SportName permanecem inalteradas — apenas sua localização estrutural muda do nível superior para o interior value.

Confluent Cloud (ksqlDB)

Passo 1 — Declare o fluxo de origem (correspondendo ao seu esquema de tópico existente):

CREATE STREAM client_events_raw (
  event_type    VARCHAR,
  user_id       VARCHAR,
  `timestamp`   VARCHAR,
  AmountWagered DOUBLE,
  BetId         VARCHAR,
  SportName     VARCHAR
) WITH (
  KAFKA_TOPIC  = 'client-raw-events',
  VALUE_FORMAT = 'JSON'
);

Passo 2 — Transforme para o envelope Xtremepush (teclas internas preservadas):

CREATE STREAM xtremepush_events
  WITH (KAFKA_TOPIC = 'xp-events', VALUE_FORMAT = 'JSON')
AS SELECT
  event_type                                         AS event,
  user_id,
  REGEXP_REPLACE(`timestamp`, 'T', ' ')             AS `timestamp`,
  STRUCT(
    AmountWagered := AmountWagered,
    BetId         := BetId,
    SportName     := SportName
  )                                                  AS value,
  STRUCT()                                           AS user_attributes
FROM client_events_raw
EMIT CHANGES;

STRUCT() constrói o objeto aninhado value com os nomes originais das propriedades do cliente. Se seus eventos têm uma variável ou um grande conjunto de propriedades, considere uma abordagem JSON pass-through usando AS_VALUE().

Amazon MSK (Kafka Streams — Java)

import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import com.fasterxml.jackson.databind.*;
import com.fasterxml.jackson.databind.node.*;
import java.util.*;
import java.util.regex.*;

ObjectMapper mapper = new ObjectMapper();
Set<String> ENVELOPE_KEYS = Set.of("event_type", "user_id", "timestamp");

StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> rawStream = builder.stream("client-raw-events");

KStream<String, String> transformed = rawStream.mapValues(raw -> {
    try {
        JsonNode src = mapper.readTree(raw);
        ObjectNode envelope = mapper.createObjectNode();

        envelope.put("event",   src.get("event_type").asText());
        envelope.put("user_id", src.get("user_id").asText());

        // Convert ISO 8601 to Xtremepush timestamp format
        String ts = src.get("timestamp").asText()
            .replace("T", " ")
            .replaceAll("\\.\\d+Z$", "");
        envelope.put("timestamp", ts);

        // Wrap original properties into value — key names preserved as-is
        ObjectNode value = mapper.createObjectNode();
        src.fields().forEachRemaining(entry -> {
            if (!ENVELOPE_KEYS.contains(entry.getKey())) {
                value.set(entry.getKey(), entry.getValue());
            }
        });
        envelope.set("value", value);
        envelope.set("user_attributes", mapper.createObjectNode());

        return mapper.writeValueAsString(envelope);
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
});

transformed.to("xp-events");

A mapValues abordagem transfere dinamicamente todos os campos que não são envelope para sem value enumerá-los explicitamente. Isso é robusto à evolução do esquema — novas propriedades que um cliente adiciona aparecerão automaticamente dentro valuede .

Google Cloud Managed Apache Kafka (Apache Beam / Dataflow)

import apache_beam as beam
import json
import re

ENVELOPE_KEYS = {'event_type', 'user_id', 'timestamp'}

class EnvelopeTransform(beam.DoFn):
    """Wraps a flat client event into the Xtremepush standard envelope.

    Inner property key names are preserved exactly as-is.
    Only structural change: properties move from top-level into value{}.
    """

    def process(self, element):
        raw = json.loads(element)

        # Convert ISO 8601 → Xtremepush timestamp format (YYYY-MM-DD HH:MM:SS)
        ts = re.sub(r'T', ' ', raw.get('timestamp', ''))
        ts = re.sub(r'\.\d+Z$', '', ts)

        # Wrap all non-envelope fields into value — key names preserved as-is
        value = {k: v for k, v in raw.items() if k not in ENVELOPE_KEYS}

        yield json.dumps({
            'event':           raw.get('event_type', ''),
            'user_id':         raw.get('user_id', ''),
            'timestamp':       ts,
            'value':           value,
            'user_attributes': {}
        })


with beam.Pipeline(options=pipeline_options) as p:
    (
        p
        | 'ReadFromKafka'   >> beam.io.ReadFromKafka(
                                   consumer_config={'bootstrap.servers': 'BROKER:9092'},
                                   topics=['client-raw-events'])
        | 'Transform'       >> beam.ParDo(EnvelopeTransform())
        | 'WriteToKafka'    >> beam.io.WriteToKafka(
                                   producer_config={'bootstrap.servers': 'BROKER:9092'},
                                   topic='xp-events')
    )