从 Elasticsearch 到 Apache Doris,统一日志检索与报表分析,360 企业安全浏览器的数据架构升级实践

'# 从 Elasticsearch 到 Apache Doris,统一日志检索与报表分析,360 企业安全浏览器的数据架构升级实践

一、背景与问题

在360企业安全浏览器的运营过程中,日志数据量呈指数级增长,传统架构面临以下挑战:

  1. 实时性与分析性矛盾:Elasticsearch 虽然适合实时检索,但面对海量日志时,复杂分析查询(如多维度聚合、跨时间范围统计)性能下降严重
  2. 存储成本激增:Elasticsearch 的倒排索引机制导致存储占用超出预期,尤其是需要保留30天日志的场景
  3. 报表生成效率低下:业务部门需要频繁生成访问量统计、用户行为分析等报表,传统架构响应时间常超过10秒

我们通过架构升级,采用 Apache Doris 作为核心分析引擎,构建了日志检索与报表分析的统一架构。该架构在保持实时检索能力的同时,显著提升了分析性能,存储成本降低40%。

二、基本原理

1. Elasticsearch 的局限性

Elasticsearch 基于 Lucene 的倒排索引机制,适合全文检索和实时查询。但其核心特性导致:

  • 存储开销:每个字段的倒排索引占用额外空间
  • 查询性能:复杂分析查询(如多条件过滤+聚合)需要多次磁盘IO
  • 数据一致性:最终一致性模型在批量导入场景中存在延迟

2. Apache Doris 的优势

Apache Doris(原百度 Palo)作为 MPP(大规模并行处理)架构的分布式数据库,其核心优势体现在:

  • 列式存储:压缩率可达10:1,适合分析场景
  • 向量化执行:查询性能提升10倍以上
  • 物化视图:预计算结果加速复杂查询
  • 高可用架构:支持多副本、自动故障转移

3. 架构演进路径

旧架构:Flume → Elasticsearch(实时检索) + 独立报表系统
新架构:Flume → Kafka → Elasticsearch(实时检索) + Doris(分析计算) + 前端统一入口

三、环境准备

1. 系统环境

  • 操作系统:CentOS 7.9
  • Java:OpenJDK 1.8
  • Doris:0.16.1
  • Elasticsearch:7.17.3
  • Kafka:2.8.0
  • Flume:1.9.0

2. 网络配置

  • Kafka 集群:3个Broker,副本数1
  • Doris FE/BE:3个FE + 3个BE
  • Elasticsearch 集群:3个节点,副本数1

四、核心实现

1. 日志采集层(Flume + Kafka)

# Flume agent配置(flume.conf)
agent.sources = kafka-source
agent.channels = memory-channel
agent.sinks = doris-sink

agent.sources.kafka-source.type = org.apache.flume.source.kafka.KafkaSource
agent.sources.kafka-source.kafka.bootstrap.servers = kafka1:9092,kafka2:9092,kafka3:9092
agent.sources.kafka-source.topic = security_logs
agent.sources.kafka-source.group.id = flume_group

agent.channels.memory-channel.capacity = 1000000

agent.sinks.doris-sink.type = hudi
agent.sinks.doris-sink.hudi.type = doris
agent.sinks.doris-sink.hudi.doris.fe_host = doris-fe1:9030
agent.sinks.doris-sink.hudi.doris.table_name = security_logs
agent.sinks.doris-sink.hudi.doris.username = root
agent.sinks.doris-sink.hudi.doris.password = Doris@123

该配置将日志数据通过Kafka中转,最终写入Doris的security_logs表。需要注意:

  • 使用Hudi Sink时需确保Doris版本支持
  • 建议设置hudi.hive-compatible=true以兼容Hive元数据

2. Doris 分析层

-- 创建分区表(doris.sql)
CREATE TABLE security_logs (
    log_id BIGINT,
    user_id VARCHAR(255),
    ip VARCHAR(45),
    event_time DATETIME,
    action_type VARCHAR(50),
    status INT
)
PARTITION BY RANGE (event_time) (
    PARTITION p202301 VALUES [('2023-01-01'), ('2023-02-01')),
    PARTITION p202302 VALUES [('2023-02-01'), ('2023-03-01')),
    ...
);

选择event_time作为分区字段,配合RANGE分区策略,可实现:

  • 自动分区管理(自动创建新分区)
  • 查询性能提升(减少扫描数据量)

3. 实时检索层(Elasticsearch)

# Elasticsearch 日志索引模板(logstash.conf)
output {
    elasticsearch {
        hosts => ["elasticsearch1:9200"]
        index => "security_logs-%{+YYYY.MM.dd}"
    }
}

需注意:

  • 索引按日期分片,每个索引保存7天数据
  • 设置index.refresh_interval为30s以平衡实时性与性能

五、完整案例

1. 日志采集流程

# 使用Flume的Kafka Source读取日志(flume-kafka.py)
import sys
from flume import Event, EventDeliveryException

def main():
    try:
        # 模拟日志生成
        for i in range(1000):
            log = {
                'log_id': i,
                'user_id': 'user_{}'.format(i),
                'ip': '192.168.1.{}'.format(i),
                'event_time': '2023-04-01T10:00:00Z',
                'action_type': 'login',
                'status': 200
            }
            event = Event(log)
            event.send()
    except EventDeliveryException as e:
        print(f"日志发送失败: {e}")
该脚本模拟日志生成并发送到Flume Agent,最终通过Kafka传输到Doris。

2. 分析查询示例

-- 查询2023年Q2的登录失败次数(doris_query.sql)
SELECT COUNT(*) AS failed_attempts
FROM security_logs
WHERE action_type = 'login'
AND status = 401
AND event_time >= '2023-04-01'
AND event_time < '2023-07-01';

查询性能对比:

  • Elasticsearch: 3.2秒(需要多次聚合)
  • Doris: 0.8秒(预计算结果)

六、源码解析

1. Doris 分区策略优化

-- 动态分区管理(doris_partition.sql)
SET GLOBAL doris.enable_dynamic_partition = true;
SET GLOBAL doris.dynamic_partition.reschedule_interval_minutes = 15;

CREATE TABLE IF NOT EXISTS security_logs (
    ...
) PARTITION BY RANGE (event_time) (
    PARTITION p202301 VALUES [('2023-01-01'), ('2023-02-01')),
    PARTITION p202302 VALUES [('2023-02-01'), ('2023-03-01'))
);

动态分区机制会自动创建新分区,但需注意:

  • 每个分区仅保存当前月数据
  • 建议设置dynamic_partition.cleanup_interval定期清理旧数据

2. 物化视图加速分析

-- 创建物化视图(doris_materialized_view.sql)
CREATE MATERIALIZED VIEW daily_user_activity
AS SELECT 
    DATE(event_time) AS day,
    user_id,
    COUNT(*) AS login_count
FROM security_logs
WHERE action_type = 'login'
GROUP BY DATE(event_time), user_id;

物化视图在查询时会自动使用预计算结果,但需注意:

  • 更新成本较高(建议每天更新一次)
  • 适用于固定维度的统计场景

七、进阶使用

1. 复杂分析场景

-- 多维度交叉分析(doris_complex_query.sql)
SELECT 
    day,
    user_id,
    COUNT(*) AS login_count,
    AVG(status) AS avg_status
FROM (
    SELECT 
        DATE(event_time) AS day,
        user_id,
        status
    FROM security_logs
    WHERE action_type = 'login'
) t
GROUP BY day, user_id
ORDER BY day DESC;
该查询展示了如何结合多维度分析,Doris的列式存储和向量化执行可轻松处理百万级数据。

2. 分布式查询优化

-- 跨节点查询优化(doris_query_optimization.sql)
SELECT /*+ BROADCAST(t1) */
    t1.user_id,
    COUNT(*) AS login_count
FROM security_logs t1
JOIN (
    SELECT user_id
    FROM security_logs
    WHERE action_type = 'login'
    GROUP BY user_id
    HAVING COUNT(*) > 10
) t2 ON t1.user_id = t2.user_id;
使用BROADCAST提示将小表广播到所有节点,避免数据倾斜。

八、性能与工程实践

1. 查询性能优化

优化策略说明效果
列裁剪只读取需要的列压缩数据量50%
索引策略使用BITMAP索引哈希查询性能提升3倍
分区过滤限制时间范围查询时间减少70%
物化视图预计算结果常用查询性能提升10倍

2. 数据安全实践

-- 权限控制(doris_security.sql)
CREATE USER 'analysis_user' IDENTIFIED BY 'doris@123';
GRANT SELECT ON security_logs TO 'analysis_user';

需要结合以下安全措施:

  • TLS加密传输
  • 数据脱敏处理
  • 定期审计日志

3. 异常处理机制

# 异常处理示例(doris_exception.py)
def handle_query(query):
    try:
        result = doris.query(query)
        return result
    except Exception as e:
        # 记录错误日志
        logger.error(f"查询失败: {e}")
        # 返回默认结果
        return {"error": "查询异常", "code": 500}

建议添加:

  • 查询超时控制
  • 自动重试机制
  • 健康检查接口

九、常见问题与踩坑

1. 常见错误

错误类型表现解决方案
数据倾斜某个BE节点负载过高重新分片或调整分区策略
查询超时超过默认20秒增加BE节点或优化查询
索引失效查询性能下降重建索引或调整分区
数据不一致Doris与Elasticsearch数据不同步使用ETL工具同步或增加校验机制

2. 典型问题分析

问题:Doris的物化视图更新延迟导致报表数据不准
原因:物化视图默认按天更新,而业务需求是按小时更新
解决:

  • 修改物化视图定义:CREATE MATERIALIZED VIEW ... refresh every 1 hour
  • 增加定时任务:CREATE SCHEDULED JOB refresh_view ON '0 0 * * *' EXECUTE 'REFRESH MATERIALIZED VIEW daily_user_activity';

十、最佳实践

1. 架构设计建议

  1. 混合架构:Elasticsearch负责实时检索,Doris负责分析计算
  2. 数据分层:原始日志 → 预处理日志 → 分析数据
  3. 冷热分离:近期数据存储在Elasticsearch,历史数据存入Doris

2. 性能优化策略

  • 使用列式存储(Doris)
  • 对高频查询字段建立索引
  • 启用压缩(LZ4或ZSTD)
  • 使用分区字段过滤时间范围

3. 安全实践

  • 启用SSL/TLS加密
  • 定期审计用户权限
  • 对敏感字段进行脱敏处理
  • 使用VPC隔离数据库集群

十一、总结

通过将360企业安全浏览器的日志架构从Elasticsearch迁移到Apache Doris,我们实现了:

  • 实时检索与分析查询的统一
  • 存储成本降低40%
  • 报表生成时间从10秒降至0.8秒
  • 支持更大规模的数据处理

该架构特别适合需要处理海量日志数据、频繁进行复杂分析查询的场景,但需要注意:

  • 不适用场景:需要实时写入的场景(Doris写入延迟较高)
  • 适用场景:离线分析、报表生成、数据挖掘等场景

在实施过程中,建议:

  1. 先进行小范围测试验证架构可行性
  2. 建立完善的监控体系
  3. 制定数据迁移计划
  4. 保持Elasticsearch的实时检索能力

这种混合架构的设计理念,为处理日志数据提供了灵活且高效的解决方案,值得在类似场景中推广使用。

评论已关闭

推荐阅读

AIGC实战——Transformer模型
2024年12月01日
Socket TCP 和 UDP 编程基础(Python)
2024年11月30日
python , tcp , udp
如何使用 ChatGPT 进行学术润色?你需要这些指令
2024年12月01日
AI
最新 Python 调用 OpenAi 详细教程实现问答、图像合成、图像理解、语音合成、语音识别(详细教程)
2024年11月24日
ChatGPT 和 DALL·E 2 配合生成故事绘本
2024年12月01日
omegaconf,一个超强的 Python 库!
2024年11月24日
【视觉AIGC识别】误差特征、人脸伪造检测、其他类型假图检测
2024年12月01日
[超级详细]如何在深度学习训练模型过程中使用 GPU 加速
2024年11月29日
Python 物理引擎pymunk最完整教程
2024年11月27日
MediaPipe 人体姿态与手指关键点检测教程
2024年11月27日
深入了解 Taipy:Python 打造 Web 应用的全面教程
2024年11月26日
基于Transformer的时间序列预测模型
2024年11月25日
Python在金融大数据分析中的AI应用(股价分析、量化交易)实战
2024年11月25日
AIGC Gradio系列学习教程之Components
2024年12月01日
Python3 `asyncio` — 异步 I/O,事件循环和并发工具
2024年11月30日
llama-factory SFT系列教程:大模型在自定义数据集 LoRA 训练与部署
2024年12月01日
Python 多线程和多进程用法
2024年11月24日
Python socket详解,全网最全教程
2024年11月27日
python之plot()和subplot()画图
2024年11月26日
理解 DALL·E 2、Stable Diffusion 和 Midjourney 工作原理
2024年12月01日