惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

D
DataBreaches.Net
N
Netflix TechBlog - Medium
F
Fortinet All Blogs
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
宝玉的分享
宝玉的分享
Y
Y Combinator Blog
博客园 - 聂微东
WordPress大学
WordPress大学
酷 壳 – CoolShell
酷 壳 – CoolShell
B
Blog RSS Feed
小众软件
小众软件
The GitHub Blog
The GitHub Blog
S
SegmentFault 最新的问题
Hugging Face - Blog
Hugging Face - Blog
Jina AI
Jina AI
Microsoft Azure Blog
Microsoft Azure Blog
V
V2EX
B
Blog
H
Help Net Security
D
Docker
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
罗磊的独立博客
月光博客
月光博客
博客园 - 司徒正美

博客园 - Hey,Coder!

微信小程序 web方案实现调音器功能 systemctl slice配置docker最大占用cpu 内存 pipenv环境变量配置 .net8验证码 数据库差异对比工具 docker部署prometheus grafana pgexporter Alertmanager ubuntu24 换源 ssh配置 ubuntu24 Harbor部署 ubuntu24 kvm部署 cockpit web管理界面 .net api代理 Ubuntu 24.04 DMAR IOMMU 报错修复 uniapp iOS App 上架 vs文件软链接 Docker 网络别名 .net8通用中间件处理公共返回值结构 .net环境下OpenTelemetry Body 读取注意事项 jaeger界面显示异常 .net opentelemetry collector jaeger c# opentelemetry 自定义export c# opentelemetry 自定义export c# 生成对象的初始化代码 portainer httphelper封装 Ocelot + Consul + SignalR 粘性会话 正交软件架构 docker postgresql17 主从复制 rabbitmq总结与c#示例 docker rabbitmq quartz dashboard docker consul registrator 服务自动注册consul 微服务动态扩容
Elasticsearch ILM 与 Data Stream 概念及实例说明文档
Hey,Coder! · 2026-07-13 · via 博客园 - Hey,Coder!

本文档基于项目中的 ESService 类代码,梳理 Elasticsearch 中 Data StreamILM / DSL 生命周期管理的核心概念、工作原理及实际代码示例。


一、Data Stream(数据流)

1. 概念

Data Stream 是 Elasticsearch 针对时序数据(日志、指标、事件)提供的一种抽象存储层。它对外表现为一个逻辑名称,底层自动管理一组后备索引(Backing Index)。用户只需向 Data Stream 写入文档,无需关心索引的创建、滚动和删除。

2. 命名规范

Data Stream 的名称必须遵循以下模式之一(用于自动匹配内置模板):

  • logs-{dataset}-{namespace} (日志)
  • metrics-{dataset}-{namespace} (指标)
  • traces-{dataset}-{namespace} (链路追踪)

本项目中使用:logs-myapp-default

3. 后备索引生成规则

当向 Data Stream 写入第一条文档时,Elasticsearch 会根据匹配的 Index Template 自动创建第一个后备索引,命名格式为:

.ds-{data_stream_name}-{yyyy.MM.dd}-{generation}

例如:.ds-logs-myapp-default-2026.07.13-000001

  • 日期:后备索引创建的日期(按天滚动)。
  • generation:从 000001 开始递增,每次 rollover 后加 1。

4. 核心优势

  • 自动滚动:配合 ILM 或 DSL,索引达到条件后自动创建新后备索引。
  • 查询优化:按时间范围查询时,ES 自动跳过不相关的后备索引(索引剪枝),提升查询速度。
  • 简化运维:无需手动管理索引生命周期。

二、生命周期管理:ILM 与 DSL

Elasticsearch 提供两种管理 Data Stream 生命周期的方式:ILM(Index Lifecycle Management)DSL(Data Stream Lifecycle)

1. ILM(传统方式,全版本支持)

概念

ILM 是一种基于索引级别的策略系统。管理员定义一个策略(Policy),其中包含多个阶段(Phase),如 hotwarmcoldfrozendelete。每个阶段可配置动作(Actions),例如 rollover、shrink、force merge、searchable snapshot、delete 等。

与 Data Stream 结合

  • 通过 Component Template 将 ILM 策略绑定到后备索引的 index.lifecycle.name setting。
  • Index Template 引用 Component Template,并标记 data_stream: {}
  • 后备索引创建后自动应用 ILM 策略,按阶段执行。

优点

  • 功能强大,支持多阶段(hot→warm→cold→frozen→delete)。
  • 支持 searchable snapshot(冷/冻阶段)。
  • 适用于自建集群和复杂场景。

缺点

  • 配置复杂(需单独建 Policy、Component Template、Index Template)。
  • 更新策略后,已有索引需手动触发 _ilm/retry 才能重新评估。

2. DSL(Data Stream Lifecycle,8.11+ 新方式)

概念

DSL 是一种基于 Data Stream 级别的简化生命周期管理。它只关心一个参数:数据保留时长(data_retention)。rollover 由 Elasticsearch 自动管理(默认 50GB 或根据保留时长自动计算滚动间隔)。

与 Data Stream 结合

  • 在 Index Template 的 template.lifecycle 中设置 data_retention
  • 也可通过 PUT _data_stream/{name}/_lifecycle 直接对已有 Data Stream 设置。
  • 无需单独的 ILM 策略,无需 Component Template。

优点

  • 配置极简,一行 DataRetention("30d") 即可。
  • 自动 rollover,无需手动设置条件。
  • 适合 Serverless 和简单场景。

缺点

  • 仅支持删除阶段,不支持 warm/cold/frozen。
  • 需要 ES 8.11+。

三、代码实例详解

项目 ESService 类提供了两种实现,分别演示 ILM 和 DSL 的用法。

1. 使用 DSL(DataStreamLifeCycle 方法)

public void DataStreamLifeCycle(int retentionDays,string logMessage)
{
    var retention = $"{retentionDays}d";
    // 1. 创建 Index Template(幂等),设置 Data Retention 为 30 分钟
    var putTemplateResponse = client.Indices.PutIndexTemplate("logs-myapp-template", it => it
        .IndexPatterns("logs-myapp-*")
        .Template(t => t
            .Lifecycle(l => l.DataRetention(retention ))   // DSL 核心:保留 30 分钟
            .Mappings(m => m
                .Properties(props => props
                    .Text("Message")
                    .Keyword("ServiceName")
                    .Keyword("LogLevel")
                    .Date("@timestamp")
                )
            )
        )
        .DataStream()
        .Priority(500)
    );

    // 2. 更新所有已存在的 Data Stream(使其立即应用新 retention)
    var listDsResponse = client.Indices.GetDataStream(ds => ds.Name("logs-myapp-*"));
    if (listDsResponse.IsValidResponse && listDsResponse.DataStreams.Any())
    {
        foreach (var dataStream in listDsResponse.DataStreams)
        {
            client.Indices.PutDataLifecycle(dataStream.Name, p => p.DataRetention(retention ));
        }
    }

    // 3. 写入一条日志,Data Stream 自动创建
    var logDoc = new LogEntry
    {
        Timestamp = DateTime.UtcNow,
        Message = logMessage,
        ServiceName = "MyService",
        LogLevel = "Info"
    };
    client.Index(logDoc, idx => idx.Index("logs-myapp-default"));
}

要点

  • 模板中 .Lifecycle(l => l.DataRetention("30m")) 定义了 DSL 保留期。
  • 对已有 Data Stream 调用 PutDataLifecycle 使其立即生效。
  • 写入时使用强类型 POCO LogEntry,并通过 [JsonPropertyName("@timestamp")] 确保字段名正确。

2. 使用 ILM(WriteLogToDataStreamAsync 方法)

public void WriteLogToDataStreamAsync(int retentionDays, string logMessage)
{
    var ilmPolicyName = $"logs-myapp-retention-{retentionDays}d";

    // 1. 创建 ILM 策略
    client.IndexLifecycleManagement.PutLifecycle(ilmPolicyName, policy => policy
        .Policy(p => p
            .Phases(ph => ph
                .Hot(h => h
                    .Actions(a => a
                        .Rollover(r => r.MaxAge("7d").MaxPrimaryShardSize("50gb"))
                    )
                )
                .Delete(d => d
                    .MinAge($"{retentionDays}d")
                    .Actions(da => da.Delete(del => del.DeleteSearchableSnapshot()))
                )
            )
        )
    );

    // 2. 创建 Component Template(绑定 ILM 策略)
    var componentTemplateName = "logs-myapp-settings";
    client.Cluster.PutComponentTemplate(componentTemplateName, ct => ct
        .Template(t => t
            .Settings(s => s
                .NumberOfShards(1)
                .NumberOfReplicas(1)
                .Lifecycle(l => l.Name(ilmPolicyName))   // 绑定 ILM 策略
            )
            .Mappings(m => m
                .Properties(props => props
                    .Text("Message")
                    .Keyword("ServiceName")
                    .Keyword("LogLevel")
                    .Date("@timestamp")
                )
            )
        )
    );

    // 3. 创建 Index Template(引用 Component Template)
    client.Indices.PutIndexTemplate("logs-myapp-template", it => it
        .IndexPatterns("logs-myapp-*")
        .ComposedOf(componentTemplateName)
        .DataStream()
        .Priority(500)
    );

    // 4. 写入日志
    var logDoc = new LogEntry { ... };
    client.Index(logDoc, idx => idx.Index("logs-myapp-default"));
}

要点

  • ILM 策略中 Hot 阶段配置 rollover(7天或50GB),Delete 阶段配置保留天数。
  • Component Template 集中管理 settings(分片、副本、ILM 策略名)和 mappings。
  • Index Template 通过 ComposedOf 引用 Component Template,并标记 DataStream()
  • 写入时自动触发 Data Stream 创建,后备索引继承 ILM 策略。

四、ILM 与 DSL 的选择建议

场景 推荐方案
只需按时间保留数据,无 warm/cold 需求 DSL(ES ≥ 8.11)
需要 warm/cold 降冷、searchable snapshot ILM
自建集群,版本较低(< 8.11) ILM
希望最小化配置,快速上线 DSL
已有大量 ILM 策略,不愿迁移 继续 ILM

注意两者不能混用。如果在 Index Template 中同时设置了 index.lifecycle.name(ILM)和 template.lifecycle(DSL),ES 默认 index.lifecycle.prefer_ilm: true,ILM 会覆盖 DSL。若要使用 DSL,必须移除 ILM 绑定。


五、常见问题与注意事项

1. 后备索引命名冲突

  • 后备索引名称由系统自动生成,不可手动修改。
  • 删除 Data Stream 时,所有后备索引会被一同删除。

2. 更新生命周期配置

  • ILM:更新策略后,已有索引需执行 POST /{index}/_ilm/retry 或等待下一次 ILM 检查周期(默认10分钟)。
  • DSL:更新 DataRetention 后,可通过 PUT _data_stream/{name}/_lifecycle 立即生效,无需额外操作。

3. 写入时必须包含 @timestamp

  • Data Stream 要求每条文档必须包含 @timestamp 字段(类型 date)。
  • 建议使用 POCO 并标记 [JsonPropertyName("@timestamp")],或使用 Dictionary<string, object> 显式设置键名。

4. 模板优先级

  • 内置模板优先级为 100,Fleet 集成为 200。
  • 自定义模板建议设置 Priority >= 300,确保覆盖内置模板。

5. 删除 Data Stream

  • 删除 Data Stream 会同时删除所有后备索引和数据,不可恢复。
  • 生产环境中应依靠生命周期自动清理,而非手动删除。

六、总结

  • Data Stream 是时序数据的最佳存储方式,配合生命周期管理可实现自动化运维。
  • ILM 功能全面但配置复杂,DSL 简洁易用但功能有限。
  • 根据 ES 版本和业务需求选择合适的方案,避免混用。
  • 代码示例展示了两种方案的完整实现,可作为项目开发的直接参考。

扩展

封装通用日志服务(同步版本,仅支持 DSL 模式),提供单条和批量创建方法

    public class EsDataStreamOptions
    {
        /// <summary>
        /// Data Stream 名称,如 "logs-myapp-default"
        /// </summary>
        public string DataStreamName { get; set; } = "logs-bll-default";

        /// <summary>
        /// 索引模板名称
        /// </summary>
        public string IndexTemplateName { get; set; } = "logs-bll-default-template";

        /// <summary>
        /// 索引匹配模式,如 "logs-default-*"
        /// DataStream 模板只匹配“后备索引(Backing Index)”的命名规则
        /// </summary>
        public string IndexPattern { get; set; } = "blllogs-default*";

        /// <summary>
        /// 数据保留天数(DSL 的 data_retention)
        /// </summary>
        public int RetentionDays { get; set; } = 30;

        /// <summary>
        /// 索引设置(分片、副本)
        /// </summary>
        public int NumberOfShards { get; set; } = 1;
        public int NumberOfReplicas { get; set; } = 1;
    }

    /// <summary>
    /// 日志基础接口,所有自定义日志类必须实现此接口
    /// </summary>
    public class BaseLogEntry
    {
        /// <summary>
        /// 
        /// </summary>
        public string Id { get; set; }

        /// <summary>
        /// 
        /// </summary>
        [JsonPropertyName("@timestamp")]
        public DateTime Timestamp { get; set; }
    }

/// <summary>
/// 通用日志服务(同步版本,仅支持 DSL 模式)
/// </summary>
public class GenericEsLogService<TLog>     where TLog : BaseLogEntry
{
    private readonly ElasticsearchClient _client;
    private EsDataStreamOptions _options;
    private readonly ILogger<GenericEsLogService<TLog>> _logger;

    public GenericEsLogService(
        ElasticsearchClient client,
        ILogger<GenericEsLogService<TLog>> logger = null)
    {
        _client = client;
        _logger = logger;
    }

    /// <summary>
    /// 
    /// </summary>
    public EsDataStreamOptions GetOption()
    {
        return _options;
    }

    /// <summary>
    /// 
    /// </summary>
    public string? GetIndexName()
    {
        return GetOption()?.DataStreamName;
    }

    /// <summary>
    /// 初始化 Data Stream 模板和生命周期(同步调用)
    /// </summary>
    public void Initialize(EsDataStreamOptions option)
    {
            this._options = option;
            //var retention = $"{_options.RetentionDays}d";
            var retention = $"{_options.RetentionDays}s";

            // 1. 创建/更新 Index Template(包含 DSL 生命周期)
            var putTemplateResponse = _client.Indices.PutIndexTemplate(
                _options.IndexTemplateName,
                it => it
                    .IndexPatterns(_options.IndexPattern)
                    .Template(t => t
                        .Lifecycle(l => l.DataRetention(retention))
                        .Mappings(m => m
                            .Properties(props =>
                            {
                                props.Date("@timestamp");
                            })
                        )
                        .Settings(s => s
                            .NumberOfShards(option.NumberOfShards)   // 只对新索引生效
                            .NumberOfReplicas(option.NumberOfReplicas)
                            .Lifecycle(f =>
                                f
                                .Name(string.Empty)
                                .PreferIlm(false)
                            )
                        )
                    )
                    //.ComposedOf("logs")
                    .DataStream()
                    .Priority(200)
            );

            if (!putTemplateResponse.IsValidResponse)
                throw new Exception($"Failed to create DSL template: {putTemplateResponse.DebugInformation}");

            _logger?.LogInformation("DSL initialization completed for {DataStream}", _options.DataStreamName);
    }


        /// <summary>
        /// 如果存在已有的数据可以更新已有数据的配置
        /// </summary>
        public void UpdateExistIndex()
        {
            // 2. 更新所有已存在的 Data Stream(使其立即应用新的 retention)
            //var listDsResponse = _client.Indices.GetDataStream(ds => ds.Name(_options.IndexPattern));
            var listDsResponse = _client.Indices.GetDataStream(ds => ds.Name(this._options.DataStreamName));
            if (listDsResponse.IsValidResponse && listDsResponse.DataStreams.Any())
            {
                foreach (var ds in listDsResponse.DataStreams)
                {
                    var putLifecycleResponse = _client.Indices.PutDataLifecycle(
                        ds.Name,
                        p => p.DataRetention(_options.RetentionDays));
                    if (!putLifecycleResponse.IsValidResponse)
                        _logger?.LogWarning("Failed to update lifecycle for DS {DsName}", ds.Name);

                    foreach (var index in ds.Indices)
                    {
                        index.PreferIlm = false;
                        var updateSettingsResponse = _client.Indices.PutSettings(index.IndexName, s => s
                            .Settings(q => q.NumberOfReplicas(_options.NumberOfReplicas)
                            .Lifecycle(l => l
                                    .Name("")                  // 空字符串 = 移除 ILM 策略绑定
                                    .PreferIlm(false)          // 关掉 ILM 优先,让 DSL 接管,保证日志正确删除
            )
                            )
                        );


                        if (!updateSettingsResponse.IsValidResponse)
                            _logger?.LogWarning("Failed to update replicas for {Index}", index.IndexName);
                    }
                }
            }
        }

    /// <summary>
    /// 写入一条日志(同步)
    /// </summary>
    public string WriteLog(TLog logEntry)
    {
        if (logEntry.Timestamp == default)
            logEntry.Timestamp = DateTime.UtcNow;

        //var indexResponse = _client.Create(logEntry, _options.DataStreamName, logEntry.Id);

        //ds只允许新增
        var indexResponse = _client.Create(logEntry, _options.DataStreamName, logEntry.Id);

        if (!indexResponse.IsValidResponse)
            throw new Exception($"Failed to write log: {indexResponse.DebugInformation}");

        _logger?.LogDebug("Log written to {DataStream}, ID={Id}", _options.DataStreamName, indexResponse.Id);

        return indexResponse.Id;
    }

    /// <summary>
    /// 批量写入日志(同步)
    /// </summary>
    public void BulkWrite(IEnumerable<TLog> logEntries)
    {
        var bulkDescriptor = new BulkRequestDescriptor();

        foreach (var entry in logEntries)
        {
            if (entry.Timestamp == default)
                entry.Timestamp = DateTime.UtcNow;

            // 新版:用 .Create(idx => idx.Index(...)) 链式追加,每个操作对应一个 Create
            bulkDescriptor.Create<TLog>(entry).Index(_options.DataStreamName);
        }

        var bulkResponse = _client.Bulk(bulkDescriptor);
        if (!bulkResponse.IsValidResponse)
            throw new Exception($"Bulk write failed: {bulkResponse.DebugInformation}");
    }
}