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 |
|---|---|---|---|
| Corda | Sim | O nome do evento, por exemplo, |
| 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. |
| Corda | Pelo menos um identificador de usuário exige | Identificador adicional de usuário, se disponível. |
| Corda | Pelo menos um identificador de usuário exige | Identificador de perfil Xtremepush. |
| Corda | Pelo menos um identificador de usuário exige | Identificador de dispositivo Xtremepush. |
| 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. |
| 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. |
| 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')
)