从 Elasticsearch 到 Apache Doris,统一日志检索与报表分析,360 企业安全浏览器的数据架构升级实践
'# 从 Elasticsearch 到 Apache Doris,统一日志检索与报表分析,360 企业安全浏览器的数据架构升级实践
一、背景与问题
在360企业安全浏览器的运营过程中,日志数据量呈指数级增长,传统架构面临以下挑战:
- 实时性与分析性矛盾:Elasticsearch 虽然适合实时检索,但面对海量日志时,复杂分析查询(如多维度聚合、跨时间范围统计)性能下降严重
- 存储成本激增:Elasticsearch 的倒排索引机制导致存储占用超出预期,尤其是需要保留30天日志的场景
- 报表生成效率低下:业务部门需要频繁生成访问量统计、用户行为分析等报表,传统架构响应时间常超过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. 架构设计建议
- 混合架构:Elasticsearch负责实时检索,Doris负责分析计算
- 数据分层:原始日志 → 预处理日志 → 分析数据
- 冷热分离:近期数据存储在Elasticsearch,历史数据存入Doris
2. 性能优化策略
- 使用列式存储(Doris)
- 对高频查询字段建立索引
- 启用压缩(LZ4或ZSTD)
- 使用分区字段过滤时间范围
3. 安全实践
- 启用SSL/TLS加密
- 定期审计用户权限
- 对敏感字段进行脱敏处理
- 使用VPC隔离数据库集群
十一、总结
通过将360企业安全浏览器的日志架构从Elasticsearch迁移到Apache Doris,我们实现了:
- 实时检索与分析查询的统一
- 存储成本降低40%
- 报表生成时间从10秒降至0.8秒
- 支持更大规模的数据处理
该架构特别适合需要处理海量日志数据、频繁进行复杂分析查询的场景,但需要注意:
- 不适用场景:需要实时写入的场景(Doris写入延迟较高)
- 适用场景:离线分析、报表生成、数据挖掘等场景
在实施过程中,建议:
- 先进行小范围测试验证架构可行性
- 建立完善的监控体系
- 制定数据迁移计划
- 保持Elasticsearch的实时检索能力
这种混合架构的设计理念,为处理日志数据提供了灵活且高效的解决方案,值得在类似场景中推广使用。
评论已关闭