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

> Discover how KCloud-Platform-IoT uses Elasticsearch for efficient log aggregation and search. Explore asynchronous indexing, custom mapping, and runtime storage switching.

- Repository: [laokou/kcloud-platform-iot](https://github.com/koushenhai/kcloud-platform-iot)
- Tags: how-to-guide
- Published: 2026-03-05

---

**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:

- **Index definition**: Declares document structure and mappings via `TraceLogIndex` in [`laokou-service/laokou-logstash/laokou-logstash-infrastructure/src/main/java/org/laokou/logstash/gatewayimpl/database/dataobject/TraceLogIndex.java`](https://github.com/koushenhai/kcloud-platform-iot/blob/main/laokou-service/laokou-logstash/laokou-logstash-infrastructure/src/main/java/org/laokou/logstash/gatewayimpl/database/dataobject/TraceLogIndex.java)
- **Storage abstraction**: Defines the `TraceLogStorage` SPI with concrete implementations for different backends
- **Configuration & wiring**: `StorageConfig` manages conditional bean creation based on the active storage type
- **Client auto-configuration**: `ElasticsearchAutoConfig` supplies pre-configured clients and templates in [`laokou-common/laokou-common-elasticsearch/src/main/java/org/laokou/common/elasticsearch/config/ElasticsearchAutoConfig.java`](https://github.com/koushenhai/kcloud-platform-iot/blob/main/laokou-common/laokou-common-elasticsearch/src/main/java/org/laokou/common/elasticsearch/config/ElasticsearchAutoConfig.java)

## 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.

```java
@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.

```java
@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.

```yaml

# 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.

```java
@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:

```yaml
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`](https://github.com/koushenhai/kcloud-platform-iot/blob/main/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.