Ir para o conteúdo

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)