Camada de Mensageria (Kafka)
Esta página documenta os componentes de mensageria responsáveis pela produção e consumo de eventos no Kafka/Redpanda.
Produtor (Producer)
O produtor realiza a serialização dos modelos de onboarding (incluindo geração automática e validação de esquemas JSON/Avro) e o envio seguro para os tópicos do Kafka.
onboarding_hub.messaging.producer
Produtor Kafka para publicar eventos de onboarding.
Implementa suporte a Avro/JSON e fallback de simulação caso o Kafka esteja indisponível.
Classes
OnboardingProducer
Produtor de mensagens de Onboarding para o Kafka.
Source code in src\onboarding_hub\messaging\producer.py
| class OnboardingProducer:
"""Produtor de mensagens de Onboarding para o Kafka."""
def __init__(self, bootstrap_servers: str | None = None) -> None:
self.bootstrap_servers = bootstrap_servers or settings.kafka_bootstrap_servers
self.topic = settings.kafka_topic_onboarding
self._producer = None
if KAFKA_AVAILABLE:
try:
conf = {
"bootstrap.servers": self.bootstrap_servers,
"client.id": "onboarding-hub-producer",
}
self._producer = Producer(conf)
logger.info("Produtor Kafka inicializado com sucesso em %s", self.bootstrap_servers)
except Exception as e:
logger.error(
"Falha ao inicializar produtor real do Kafka: %s. Operando em modo SIMULAÇÃO.",
e,
)
def publish_onboarding(self, onboarding: OnboardingPendente) -> None:
"""Publica um evento de onboarding pendente."""
payload_dict = json.loads(onboarding.model_dump_json())
if self._producer is not None:
try:
# Callback de entrega
def delivery_report(err: Any, msg: Any) -> None:
if err is not None:
logger.error("Falha ao entregar mensagem no Kafka: %s", err)
else:
logger.debug("Mensagem entregue em %s [%d]", msg.topic(), msg.partition())
# Serializa para JSON (em produção real, usaria-se o AvroSerializer
# com Schema Registry)
payload_bytes = json.dumps(payload_dict).encode("utf-8")
self._producer.produce(
topic=self.topic,
key=str(onboarding.id_cadastro),
value=payload_bytes,
callback=delivery_report,
)
# Dispara callbacks em background
self._producer.poll(0)
logger.info(
"Evento %s publicado no Kafka (tópico: %s)", onboarding.id_cadastro, self.topic
)
except Exception as e:
logger.error("Erro ao enviar para o Kafka: %s. Salvando mock log.", e)
self._mock_log(payload_dict)
else:
self._mock_log(payload_dict)
def flush(self, timeout: float = 1.0) -> int:
"""Garante que todas as mensagens pendentes sejam enviadas."""
if self._producer is not None:
return self._producer.flush(timeout)
return 0
def _mock_log(self, payload: dict[str, Any]) -> None:
"""Grava log de simulação em modo offline."""
logger.info(
"[SIMULAÇÃO KAFKA] Evento publicado no tópico '%s' para cadastro ID %s: %s",
self.topic,
payload.get("id_cadastro"),
json.dumps(payload)[:120] + "...",
)
print(f"[SIMULAÇÃO KAFKA] ID: {payload.get('id_cadastro')} -> Enviado.")
|
Methods:
flush(timeout=1.0)
Garante que todas as mensagens pendentes sejam enviadas.
Source code in src\onboarding_hub\messaging\producer.py
| def flush(self, timeout: float = 1.0) -> int:
"""Garante que todas as mensagens pendentes sejam enviadas."""
if self._producer is not None:
return self._producer.flush(timeout)
return 0
|
publish_onboarding(onboarding)
Publica um evento de onboarding pendente.
Source code in src\onboarding_hub\messaging\producer.py
| def publish_onboarding(self, onboarding: OnboardingPendente) -> None:
"""Publica um evento de onboarding pendente."""
payload_dict = json.loads(onboarding.model_dump_json())
if self._producer is not None:
try:
# Callback de entrega
def delivery_report(err: Any, msg: Any) -> None:
if err is not None:
logger.error("Falha ao entregar mensagem no Kafka: %s", err)
else:
logger.debug("Mensagem entregue em %s [%d]", msg.topic(), msg.partition())
# Serializa para JSON (em produção real, usaria-se o AvroSerializer
# com Schema Registry)
payload_bytes = json.dumps(payload_dict).encode("utf-8")
self._producer.produce(
topic=self.topic,
key=str(onboarding.id_cadastro),
value=payload_bytes,
callback=delivery_report,
)
# Dispara callbacks em background
self._producer.poll(0)
logger.info(
"Evento %s publicado no Kafka (tópico: %s)", onboarding.id_cadastro, self.topic
)
except Exception as e:
logger.error("Erro ao enviar para o Kafka: %s. Salvando mock log.", e)
self._mock_log(payload_dict)
else:
self._mock_log(payload_dict)
|
Consumidor (Consumer)
O consumidor escuta os tópicos e processa as mensagens recebidas de maneira assíncrona, efetuando as regras de negócio de onboarding e inserindo os dados no banco relacional.
onboarding_hub.messaging.consumer
Consumidor Kafka para processar cadastros de onboarding pendentes.
Classes
OnboardingConsumer
Consumidor de mensagens de Onboarding do Kafka.
Source code in src\onboarding_hub\messaging\consumer.py
| class OnboardingConsumer:
"""Consumidor de mensagens de Onboarding do Kafka."""
def __init__(self, bootstrap_servers: str | None = None, group_id: str | None = None) -> None:
self.bootstrap_servers = bootstrap_servers or settings.kafka_bootstrap_servers
self.group_id = group_id or settings.kafka_group_id
self.topic = settings.kafka_topic_onboarding
self._consumer = None
self._running = False
if KAFKA_AVAILABLE:
try:
conf = {
"bootstrap.servers": self.bootstrap_servers,
"group.id": self.group_id,
"auto.offset.reset": "earliest",
"enable.auto.commit": True,
}
self._consumer = Consumer(conf)
self._consumer.subscribe([self.topic])
logger.info(
"Consumidor Kafka subscrito no tópico '%s' (grupo: %s)",
self.topic,
self.group_id,
)
except Exception as e:
logger.error(
"Falha ao inicializar consumidor real do Kafka: %s. "
"Operando em modo SIMULAÇÃO.",
e,
)
def consume(
self, limit: int = 10, timeout: float = 1.0
) -> Generator[OnboardingPendente, None, None]:
"""Consome mensagens do Kafka e retorna instâncias de OnboardingPendente."""
if self._consumer is not None:
count = 0
while count < limit:
msg = self._consumer.poll(timeout)
if msg is None:
break
err = msg.error()
if err is not None:
if err.code() != KafkaError._PARTITION_EOF:
logger.error("Erro no consumidor Kafka: %s", err)
break
val = msg.value()
if val is None:
continue
try:
payload = json.loads(val.decode("utf-8"))
yield OnboardingPendente.model_validate(payload)
count += 1
except Exception as e:
logger.error("Falha ao decodificar mensagem do Kafka: %s", e)
else:
# Simulação: gera registros sintéticos se estiver rodando em modo offline
from onboarding_hub.synthetic.generator import OnboardingSyntheticGenerator
logger.info(
"[SIMULAÇÃO KAFKA] Gerando registros falsos para simular consumo do tópico '%s'",
self.topic,
)
generator = OnboardingSyntheticGenerator()
yield from generator.generate_records(limit)
def close(self) -> None:
"""Fecha a conexão do consumidor."""
if self._consumer is not None:
self._consumer.close()
logger.info("Consumidor Kafka fechado.")
|
Methods:
close()
Fecha a conexão do consumidor.
Source code in src\onboarding_hub\messaging\consumer.py
| def close(self) -> None:
"""Fecha a conexão do consumidor."""
if self._consumer is not None:
self._consumer.close()
logger.info("Consumidor Kafka fechado.")
|
consume(limit=10, timeout=1.0)
Consome mensagens do Kafka e retorna instâncias de OnboardingPendente.
Source code in src\onboarding_hub\messaging\consumer.py
| def consume(
self, limit: int = 10, timeout: float = 1.0
) -> Generator[OnboardingPendente, None, None]:
"""Consome mensagens do Kafka e retorna instâncias de OnboardingPendente."""
if self._consumer is not None:
count = 0
while count < limit:
msg = self._consumer.poll(timeout)
if msg is None:
break
err = msg.error()
if err is not None:
if err.code() != KafkaError._PARTITION_EOF:
logger.error("Erro no consumidor Kafka: %s", err)
break
val = msg.value()
if val is None:
continue
try:
payload = json.loads(val.decode("utf-8"))
yield OnboardingPendente.model_validate(payload)
count += 1
except Exception as e:
logger.error("Falha ao decodificar mensagem do Kafka: %s", e)
else:
# Simulação: gera registros sintéticos se estiver rodando em modo offline
from onboarding_hub.synthetic.generator import OnboardingSyntheticGenerator
logger.info(
"[SIMULAÇÃO KAFKA] Gerando registros falsos para simular consumo do tópico '%s'",
self.topic,
)
generator = OnboardingSyntheticGenerator()
yield from generator.generate_records(limit)
|