How KCloud-Platform-IoT Uses Elasticsearch for Log Aggregation and Search

KCloud-Platform-IoT leverages Elasticsearch as the default storage engine for centralized log aggregation, providing asynchronous bulk indexing, custom mapping annotations, and a pluggable storage SPI that allows runtime switching between Elasticsearch and Loki without code changes.

The KCloud-Platform-IoT project implements a robust, production-ready logging pipeline using Elasticsearch for log aggregation and search capabilities. This open-source IoT platform stores trace logs in Elasticsearch indices with near real-time refresh intervals, enabling fast full-text search across distributed services. The architecture decouples storage concerns through a clean SPI (Service Provider Interface) design, making it straightforward to configure, extend, or swap underlying storage technologies.

Architecture Overview

The logging infrastructure is organized into four distinct layers, each with specific responsibilities:

Defining the Index Structure with Custom Annotations

The platform uses custom Laokou annotations to define Elasticsearch mappings declaratively. The TraceLogIndex class specifies field types, formats, and index settings, including a refresh_interval of 1s for near real-time searchability.

@Index(setting = @Setting(refreshInterval = "1"))
public final class TraceLogIndex implements Serializable {
    @Field(type = Type.KEYWORD)
    private String serviceId;
    
    @Field(type = Type.KEYWORD)
    private String profile;
    
    @Field(type = Type.DATE, format = DateConstants.YYYY_B_MM_B_DD_HH_R_MM_R_SS_D_SSS)
    private String dateTime;
    
    @Field(type = Type.KEYWORD, index = true)
    private String traceId;
    // additional fields...
}

These annotations automatically generate the appropriate Elasticsearch mapping JSON when the index is created.

Implementing the Storage Layer

The TraceLogStorage SPI

The TraceLogStorage interface defines the contract for log persistence, abstracting whether the underlying storage is Elasticsearch, Loki, or another backend. This allows the log ingestion components to remain agnostic of the specific storage implementation.

TraceLogElasticsearchStorage Implementation

The TraceLogElasticsearchStorage class implements this interface and handles the actual interaction with Elasticsearch. Its batchSave method performs three critical operations: ensuring the index exists, bulk-inserting documents, and blocking until completion using virtual threads.

@Override
public void batchSave(List<Object> list) {
    elasticsearchIndexTemplate
        .asyncCreateIndex(getIndexName(), TRACE_INDEX, TraceLogIndex.class, virtualThreadExecutor)
        .thenComposeAsync(
            _ -> elasticsearchDocumentTemplate.asyncBulkCreateDocuments(getIndexName(), list, virtualThreadExecutor),
            virtualThreadExecutor)
        .join();
}

This implementation relies on beans provided by ElasticsearchAutoConfig: ElasticsearchDocumentTemplate for document CRUD operations, ElasticsearchIndexTemplate for index lifecycle management, and a virtual thread ExecutorService for non-blocking async execution.

Conditional Configuration and Backend Selection

The StorageConfig class uses Spring's @ConditionalOnProperty annotation to select the active storage implementation at runtime based on the storage.type property.


# application.yml

storage:
  type: elasticsearch  # Change to 'loki' to switch backends

When storage.type=ELASTICSEARCH, the configuration class instantiates TraceLogElasticsearchStorage. Changing this property to loki automatically wires the Loki implementation instead, requiring no code modifications or recompilation.

Auto-Configuring the Elasticsearch Client

The ElasticsearchAutoConfig class bootstraps the Elasticsearch Java client (version 9.3.1) and wraps it in high-level templates. This auto-configuration creates both synchronous and asynchronous clients (ElasticsearchClient and ElasticsearchAsyncClient) along with three specialized templates:

  • ElasticsearchDocumentTemplate: Handles single and bulk document operations
  • ElasticsearchIndexTemplate: Manages index creation, deletion, and mapping updates
  • ElasticsearchSearchTemplate: Provides a fluent DSL for query execution

These beans are conditionally created when SpringElasticsearchProperties is present on the classpath, ensuring graceful degradation if Elasticsearch dependencies are excluded.

End-to-End Log Ingestion Flow

The complete log aggregation process follows this sequence:

  1. Application startup: Spring loads ElasticsearchAutoConfig, creating client and template beans
  2. Storage initialization: StorageConfig reads storage.type and injects the appropriate TraceLogStorage implementation
  3. Log collection: A log collector service obtains the TraceLogStorage bean via dependency injection
  4. Persistence: The service calls batchSave(messages), which converts raw logs to TraceLogIndex objects
  5. Indexing: TraceLogElasticsearchStorage ensures the index exists and executes a bulk insert operation
  6. Availability: Documents become searchable within one second due to the refresh_interval setting

Searching and Querying Trace Logs

Querying logs utilizes the ElasticsearchSearchTemplate, which abstracts the low-level Elasticsearch Java client API. The template returns Spring-style Page<T> results for easy integration with web controllers.

@Autowired
private ElasticsearchSearchTemplate searchTemplate;

public Page<TraceLogIndex> searchByTraceId(String traceId, int page, int size) {
    Query query = MultiMatchQuery.of(m -> m
            .fields("traceId")
            .query(traceId))
            ._toQuery();

    return searchTemplate.search(
            List.of("trace_log_index"),
            query,
            TraceLogIndex.class,
            PageRequest.of(page, size));
}

This approach supports complex queries while maintaining type safety and avoiding raw JSON construction.

Switching Between Elasticsearch and Loki

The pluggable architecture enables switching storage backends by modifying a single configuration value. To use Loki instead of Elasticsearch for log aggregation and search:

storage:
  type: loki

The StorageConfig class detects this change and instantiates TraceLogLokiStorage rather than TraceLogElasticsearchStorage. All upstream components continue to function normally because they depend only on the TraceLogStorage interface, not concrete implementations.

Summary

  • KCloud-Platform-IoT implements a layered architecture separating index definitions (TraceLogIndex), storage implementations (TraceLogElasticsearchStorage), and client configuration (ElasticsearchAutoConfig).
  • The TraceLogIndex class in laokou-service/laokou-logstash/laokou-logstash-infrastructure uses custom annotations to define mappings with a 1-second refresh interval for near real-time search.
  • Asynchronous bulk indexing uses virtual threads to achieve high throughput without blocking the log collection pipeline.
  • Spring's @ConditionalOnProperty in StorageConfig enables zero-code switching between Elasticsearch and Loki backends via the storage.type property.
  • High-level templates (ElasticsearchDocumentTemplate, ElasticsearchSearchTemplate) simplify CRUD operations and complex queries while maintaining full access to the underlying Elasticsearch Java client 9.3.1.

Frequently Asked Questions

How does KCloud-Platform-IoT handle high-volume log ingestion?

The platform handles high throughput via TraceLogElasticsearchStorage, which uses asyncBulkCreateDocuments combined with a virtual thread ExecutorService to perform non-blocking bulk inserts. The batchSave method chains index creation and document insertion asynchronously, blocking only at the end of the operation to ensure durability without consuming platform threads during I/O wait.

Can I use a different storage backend without modifying the source code?

Yes. By changing the storage.type property in application.yml from elasticsearch to loki, the StorageConfig class automatically wires the TraceLogLokiStorage bean instead of TraceLogElasticsearchStorage. This design follows the Strategy pattern, allowing runtime backend selection without recompilation or code changes.

What Elasticsearch client version does KCloud-Platform-IoT use?

The platform uses the official Elastic Java client version 9.3.1, as configured in ElasticsearchAutoConfig. This provides both synchronous (ElasticsearchClient) and asynchronous (ElasticsearchAsyncClient) clients, along with high-level templates that abstract connection management, retry logic, and request building.

How are Elasticsearch indices created automatically?

The TraceLogElasticsearchStorage class checks for index existence before each bulk operation and calls asyncCreateIndex via ElasticsearchIndexTemplate if the index does not exist. This method uses the TraceLogIndex class annotations to generate the correct mappings and settings, ensuring the index is properly configured before documents are inserted.

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 →