Using RabbitMQ and Kafka for Asynchronous Communication in KCloud-Platform-IoT

KCloud-Platform-IoT leverages Apache Kafka as the primary message broker for high-throughput asynchronous communication, while architecturally supporting RabbitMQ as an optional alternative for traditional queue-based workloads.

KCloud-Platform-IoT implements an event-driven microservices architecture that relies on asynchronous messaging to decouple services and ensure reliable data flow across distributed components. The platform provides out-of-the-box Kafka integration with centralized topic definitions, custom binary serialization, and domain-driven event handlers. While RabbitMQ is not fully implemented in the current codebase, the Spring-based architecture anticipates its addition through interchangeable templates and listener abstractions.

Architecture Overview

The messaging layer in KCloud-Platform-IoT is designed around event-driven communication patterns that prioritize durability and ordered delivery. The system abstracts broker-specific details behind Spring interfaces, allowing operators to choose between streaming (Kafka) and queuing (RabbitMQ) semantics based on throughput and ordering requirements.

Kafka as the Primary Event Broker

Apache Kafka serves as the default backbone for asynchronous communication in KCloud-Platform-IoT. The implementation uses custom serializers optimized for compact binary payloads and centralized topic enumeration via Java enums. According to the source code in laokou-common/laokou-common-kafka, the platform serializes objects using the ForyFactory before transmission, ensuring efficient network utilization.

RabbitMQ Extensibility

Although the current codebase does not contain active RabbitMQ configurations, the architecture deliberately anticipates RabbitMQ integration. Developers can introduce RabbitMQ by creating a RabbitConfig bean and utilizing @RabbitListener annotations alongside existing Kafka consumers. This design allows both brokers to coexist without code duplication, supporting hybrid deployments where Kafka handles high-throughput telemetry and RabbitMQ manages classic work queues.

Implementing Kafka in KCloud-Platform-IoT

The Kafka integration follows a structured approach: centralized topic definition, programmatic creation, custom serialization, and domain-specific consumption.

Centralized Topic Management with MqEnum

All logical messaging channels are enumerated in laokou-service/laokou-oss/laokou-oss-domain/src/main/java/org/laokou/oss/model/MqEnum.java. This enum defines topic names, partition counts, and replication factors in a single source of truth.

public enum MqEnum {
    OSS_LOG_TOPIC("oss-trace-log", 3, (short) 1),
    GATEWAY_TRACE_LOG_TOPIC("gateway-trace-log", 3, (short) 1),
    OPERATE_LOG_TOPIC("operate-log", 3, (short) 1);
    
    private final String topic;
    private final int partitions;
    private final short replication;
    
    // constructors and getters
}

Programmatic Topic Creation

Topics are created automatically on application startup via Log4j2Config in laokou-common/laokou-common-log4j2/src/main/java/org/laokou/common/log4j2/config/Log4j2Config.java. This configuration class injects a KafkaAdmin bean and declares NewTopic instances for each enumerated channel.

@Configuration
public class Log4j2Config {
    @Bean
    public KafkaAdmin kafkaAdmin(KafkaProperties properties) {
        return new KafkaAdmin(properties.buildAdminProperties());
    }

    @Bean
    public KafkaAdmin.NewTopics newTopics() {
        return new KafkaAdmin.NewTopics(
            new NewTopic(MqEnum.GATEWAY_TRACE_LOG_TOPIC.getTopic(), 
                        3, (short) 1),
            new NewTopic(MqEnum.OSS_LOG_TOPIC.getTopic(), 
                        3, (short) 1)
        );
    }
}

Binary Serialization with ForyFactory

The platform uses a custom serialization strategy to minimize payload size. Located in laokou-common/laokou-common-kafka/src/main/java/org/laokou/common/kafka/config/, the ForyKafkaSerializer delegates to ForyFactory.INSTANCE for compact binary encoding.

public final class ForyKafkaSerializer implements Serializer<Object> {
    @Override
    public byte[] serialize(String topic, Object data) {
        return ForyFactory.INSTANCE.serialize(data);
    }
}

Symmetric deserialization is handled by ForyKafkaDeserializer, which restores the original Java object from the byte array before delivery to consumers.

Producing Events

Any service can publish events by injecting Spring's KafkaTemplate. The template is configured to use the custom ForyKafkaSerializer, allowing direct transmission of domain objects without manual conversion.

@Service
public class TraceLogPublisher {
    private final KafkaTemplate<String, Object> kafkaTemplate;
    private static final String TOPIC = MqEnum.GATEWAY_TRACE_LOG_TOPIC.getTopic();

    public TraceLogPublisher(KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void publish(TraceLog log) {
        kafkaTemplate.send(TOPIC, log);
    }
}

Consuming with Domain Handlers

Consumers are implemented as domain event handlers using the @KafkaListener annotation. The OperateEventHandler in laokou-common/laokou-common-log/src/main/java/org/laokou/common/log/handler/OperateEventHandler.java demonstrates this pattern by listening to the operate-log topic and processing events for persistence or auditing.

@Component
public class OperateEventHandler {
    @KafkaListener(
        topics = MqEnum.OPERATE_LOG_TOPIC.getTopic(), 
        groupId = "operate-log-consumer-group"
    )
    public void handle(OperateLog log) {
        // Persist to database or trigger downstream workflows
    }
}

RabbitMQ Integration Strategy

While Kafka handles the current workload, the architecture supports RabbitMQ through analogous Spring abstractions.

Configuration Pattern

Adding RabbitMQ requires defining a RabbitConfig class that exposes ConnectionFactory, RabbitTemplate, and SimpleMessageListenerContainer beans. This configuration would mirror the Kafka setup, using a hypothetical RabbitQueueEnum to centralize queue definitions.

Producer and Consumer Implementation

RabbitMQ producers use RabbitTemplate.convertAndSend() instead of KafkaTemplate.send(). Consumers replace @KafkaListener with @RabbitListener, maintaining identical business logic within the handler methods.

@Service
public class RabbitTraceLogPublisher {
    private final RabbitTemplate rabbitTemplate;
    private static final String QUEUE = "gateway-trace-queue";

    public RabbitTraceLogPublisher(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    public void publish(TraceLog log) {
        rabbitTemplate.convertAndSend(QUEUE, log);
    }
}
@Component
public class RabbitTraceLogConsumer {
    @RabbitListener(queues = "gateway-trace-queue")
    public void onMessage(TraceLog log) {
        // Identical processing to Kafka consumer
    }
}

Summary

  • KCloud-Platform-IoT uses Kafka as the primary broker for high-throughput, ordered event streams, with implementations found in the laokou-common-kafka module.
  • Topic definitions are centralized in MqEnum.java, ensuring consistent naming and configuration across producers and consumers.
  • Custom serialization via ForyKafkaSerializer and ForyKafkaDeserializer optimizes payload size using the ForyFactory binary protocol.
  • Domain handlers consume events using @KafkaListener annotations, processing logs and telemetry in OperateEventHandler and similar components.
  • RabbitMQ is architecturally supported as an optional replacement, with clear extension points for RabbitTemplate and @RabbitListener implementations.

Frequently Asked Questions

Does KCloud-Platform-IoT support RabbitMQ out of the box?

No, the current implementation primarily supports Apache Kafka with full serialization and topic management implementations. However, the Spring-based architecture is designed to accommodate RabbitMQ as a drop-in alternative by adding the appropriate configuration beans and listener annotations without modifying core business logic.

How are Kafka topics defined and created in the codebase?

Topics are defined as enum constants in MqEnum.java, which specifies the topic name, partition count, and replication factor. These definitions are referenced by Log4j2Config.java, which programmatically creates the topics during application startup using Spring Kafka's KafkaAdmin and NewTopic APIs.

What serialization method does KCloud-Platform-IoT use for Kafka messages?

The platform uses a custom binary serialization strategy implemented in ForyKafkaSerializer.java. This serializer delegates to ForyFactory.INSTANCE.serialize(), producing compact byte arrays that minimize network overhead compared to JSON or XML formats.

Can Kafka and RabbitMQ coexist in the same KCloud-Platform-IoT deployment?

Yes, the architecture allows both brokers to coexist simultaneously. Services can inject both KafkaTemplate and RabbitTemplate, routing high-throughput telemetry through Kafka while using RabbitMQ for traditional work queues or complex routing scenarios. This hybrid approach leverages Spring's messaging abstractions to prevent code duplication between the two systems.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →