实时监控DNS敏感域名全方案分析与落地实现
一、需求完整拆解、痛点与拓展问题
1. 核心业务需求
- 数据源:多台DNS服务器本地实时日志文件,每条日志包含客户端IP、查询域名;日志持续滚动写入,高吞吐实时流。
- 匹配规则:敏感域名库存储在星环Inceptor,千万级存量、支持实时新增/删除敏感域名;每条DNS日志需和敏感库做匹配。
- 业务逻辑:日志域名命中敏感库 → 生成告警(客户端IP、命中域名、DNS服务器、时间戳、日志原文)。
- 持久化:告警数据落地存储(支持查询、溯源、统计)。
- 时效要求:实时监控,延迟控制在秒级。
2. 核心技术痛点(原生难点)
- Inceptor千万级大表,每条日志直接JDBC查表匹配,IO、网络、计算开销爆炸,完全无法实时;
- DNS日志多服务器分散、日志切割滚动、日志量大(峰值十万/秒),文件采集丢数、延迟风险高;
- 敏感域名实时更新,缓存需动态同步,不能一次性全量加载后永久不变;
- 流处理高并发,匹配、去重、窗口聚合、告警持久化需解耦,避免单点阻塞;
- 告警数据长期存储,需兼顾明细查询+报表统计,冷热数据分层;
- 分布式多DNS节点,统一采集、统一规则匹配、统一告警输出。
3. 拓展衍生需求(生产必考虑)
- 日志容错:日志轮转、文件删除、断流、重复日志去重;
- 敏感库更新同步:Inceptor表增删敏感域名,分钟级同步到内存匹配引擎;
- 告警降噪:同一IP短时间高频查询同一敏感域名,批量聚合告警,防止风暴;
- 分级告警:区分高危/普通敏感域名,支持不同推送渠道(短信、钉钉、邮件);
- 溯源能力:原始日志留存、告警支持按IP/域名/时间检索;
- 横向扩容:DNS节点增加、流量翻倍时架构线性扩容;
- 监控运维:采集延迟、匹配失败率、敏感库同步延迟、系统负载监控;
- 权限审计:敏感域名库变更操作日志留存,告警数据只读审计;
- 黑名单白名单:支持白名单IP/白名单域名,命中后不生成告警。
二、整体技术架构分层(四层架构)
架构总图
DNS服务器集群(多节点日志文件)
↓ 采集层
Filebeat 日志采集器(单机部署)
↓ 消息缓冲层
Kafka集群(日志Topic、敏感域名同步Topic、告警Topic)
↓ 实时计算匹配层
Flink流式计算引擎(核心匹配服务)
├─ Inceptor定时/增量拉取敏感黑名单 → 本地内存缓存(Redis兜底)
├─ 流解析、清洗、域名匹配、告警聚合降噪
↓ 存储层
1. 告警明细库:Inceptor(离线统计、全量溯源)
2. 实时告警缓存:Redis(告警聚合、降噪、快速查询)
3. 原始日志归档:HDFS(冷数据长期存储)
4. 运维监控:Prometheus + Grafana
各层职责与选型原因
1. 采集层:Filebeat
- 作用:多台DNS服务器本地采集滚动日志,轻量无侵入,不占用DNS服务资源;
- 优势:自带文件断点续传、日志切割适配、低资源消耗、支持过滤清洗;
- 排除Logstash:JVM重,DNS服务器资源紧张场景不适合,Filebeat轻量化。
2. 消息缓冲层:Kafka
- 作用:削峰填谷,隔离采集与计算;多分区并行消费,支撑高并发DNS流量;
- Topic划分:
- dns-log-raw:原始DNS查询日志(采集上报)
- sensitive-domain-sync:敏感域名增量变更同步流
- dns-alarm-data:生成后的告警消息
- 选型理由:高吞吐、持久化、可回溯,Flink原生完美对接;应对流量突增避免计算层崩溃。
3. 实时计算层:Apache Flink(核心)
核心能力:
- 为什么不用Spark Streaming:Spark微批分钟级延迟,无法满足实时秒级监控;Flink原生流处理,低延迟。
4. 缓存层:Redis(配套Flink)
- 用途1:存储全量敏感域名HashSet,Flink本地缓存失效兜底;
- 用途2:告警降噪去重,窗口内重复命中IP+域名计数;
- 用途3:缓存最新一批敏感域名,加速增量同步;
5. 存储层
- 两张核心表:敏感域名黑名单表、DNS告警明细表;
- 优势:兼容Hive SQL,海量结构化数据离线统计、多维度查询;千万级黑名单可增量抽取;
6. 运维监控层:Prometheus + Grafana
监控指标:日志采集延迟、Kafka堆积量、Flink任务延迟、敏感库同步耗时、告警生成量、匹配失败率。
三、资源配置方案(分流量档位,生产标准)
场景定义
峰值流量:单秒DNS查询日志 5万条;敏感域名表1000万条,每5分钟增量更新。
1. Filebeat(DNS服务器侧)
- 部署:每台DNS服务器单机部署1实例;
- 资源:CPU 0.2核,内存256M;磁盘占用极低;
- 配置:本地文件断点记录,压缩发送Kafka。
2. Kafka集群(3节点)
- 单节点配置:8核16G,1T机械磁盘;
- Topic分区:dns-log-raw 16分区,dns-alarm-data 8分区;
- 副本数2,保证消息不丢失。
3. Flink实时计算集群(Standalone集群)
- JobManager:1台 4核8G;
- TaskManager:4台 8核16G,每台slot=4;总并发16,支撑5w/s日志处理;
- Checkpoint存储:HDFS,间隔30s;
- 内存优化:堆内存12G,托管内存2G,适配千万域名内存缓存。
4. Redis集群(3主3从)
- 单节点:4核8G,SSD;用于敏感域名缓存、告警去重;
5. Inceptor集群(星环现有集群复用)
- 已有数仓集群,无需额外新增;仅新增2张业务表;
- 黑名单表分区:按更新日期分区,加速增量拉取;告警表按天分区存储。
6. HDFS
复用星环配套HDFS,存储原始日志归档、Flink Checkpoint。
四、技术路线、原理与设计思路
1. 核心难点:千万级敏感域名实时匹配方案(关键)
方案原理:分层缓存,避免每次查表Inceptor
- 定时任务:每小时全量拉取Inceptor敏感域名表,构建布隆过滤器+HashMap存入Flink TaskManager本地内存;
- 增量同步:业务端更新Inceptor敏感域名时,同步将新增/删除域名发送至Kafka sensitive-domain-sync;Flink消费增量流,实时更新本地内存缓存;
1)布隆过滤器快速预过滤:不存在直接跳过,减少精确匹配;
2)HashMap精确匹配:布隆命中后校验,消除误判;
禁止方案:每条DNS日志直接JDBC查询Inceptor,QPS上万时数据库直接打垮,延迟分钟级,不满足实时。
2. 完整数据流处理步骤
- 步骤1:日志清洗解析,提取client_ip、query_domain、dns_host、log_time,过滤非法空域名;
- 步骤2:读取本地内存敏感域名缓存,执行布隆+Hash双层匹配;
- 步骤3:命中敏感域名 → 封装告警实体(IP、域名、服务器、时间、日志原文、敏感域名等级);
- 步骤4:5分钟滚动窗口聚合,同一IP+同一域名合并计数,实现告警降噪,避免刷屏;
- 步骤5:聚合后告警写入Kafka dns-alarm-data;
3. 容错与一致性设计
- Filebeat:offset持久化,日志轮转/服务器重启不丢日志;
- Kafka:消息持久化,副本机制,消费offset提交保证至少一次消费;
- Flink Checkpoint:30s一次快照,任务崩溃重启从快照恢复,重复日志通过IP+时间戳去重;
- Inceptor写入:批量写入+事务,避免告警重复写入;
- 敏感缓存一致性:增量流+定时全量双重保障,防止缓存和Inceptor黑名单不一致。
五、数据表设计(Inceptor)
表1:敏感域名黑名单表 ods_sensitive_domain_blacklist(千万级,实时更新)
CREATE TABLE ods_sensitive_domain_blacklist (
domain STRING COMMENT '敏感域名,如xxx.com',
risk_level TINYINT COMMENT '风险等级 1高危 2普通',
create_time TIMESTAMP COMMENT '录入时间',
update_time TIMESTAMP COMMENT '更新时间',
is_delete TINYINT COMMENT '0有效 1删除'
)
PARTITIONED BY (dt STRING COMMENT '日期分区 yyyyMMdd')
STORED AS ORC
TBLPROPERTIES ('orc.compress'='snappy');
表2:DNS敏感告警明细表 dwd_dns_query_alarm(告警存储)
CREATE TABLE dwd_dns_query_alarm (
client_ip STRING COMMENT '查询客户端IP',
sensitive_domain STRING COMMENT '命中的敏感域名',
risk_level TINYINT COMMENT '风险等级',
dns_server_host STRING COMMENT '产生日志的DNS服务器',
query_time TIMESTAMP COMMENT 'DNS原始查询时间',
alarm_count BIGINT COMMENT '窗口内命中次数',
raw_log STRING COMMENT '原始日志原文',
sync_time TIMESTAMP COMMENT '告警入库时间'
)
PARTITIONED BY (dt STRING, hour STRING)
STORED AS ORC
TBLPROPERTIES ('orc.compress'='snappy');
六、完整落地代码实现
前置依赖
Flink 1.17、Flink Kafka Connector、Hive/Inceptor JDBC、Redis客户端、布隆过滤器工具、Scala/Java(下面使用Java实现主流方案)
Maven核心依赖:
<!– Flink核心 –>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.17.0</version>
</dependency>
<!– Inceptor/Hive JDBC –>
<dependency>
<groupId>org.apache.hive</groupId>
<artifactId>hive-jdbc</artifactId>
<version>3.1.2</version>
</dependency>
<!– Redis –>
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>4.4.3</version>
</dependency>
<!– 布隆过滤器 –>
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>32.1.3-jre</version>
</dependency>
代码1:DNS日志实体、告警实体
import lombok.Data;
import java.sql.Timestamp;
/** 原始DNS日志解析实体 */
@Data
public class DnsLog {
private String clientIp;
private String queryDomain;
private String dnsHost;
private Timestamp queryTime;
private String rawLog;
}
/** 告警实体 */
@Data
public class DnsAlarm {
private String clientIp;
private String sensitiveDomain;
private Integer riskLevel;
private String dnsServerHost;
private Timestamp queryTime;
private Long alarmCount;
private String rawLog;
private Timestamp syncTime;
}
代码2:敏感域名缓存管理器(核心,布隆+HashMap,增量同步)
import com.google.common.hash.BloomFilter;
import com.google.common.hash.Funnels;
import redis.clients.jedis.JedisCluster;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
public class SensitiveDomainCacheManager {
// 千万级布隆过滤器,预期插入1000万,误判率0.001
private static final long EXPECTED_INSERTIONS = 10_000_000L;
private static final double FPP = 0.001;
private BloomFilter<String> bloomFilter;
private Map<String, Integer> domainRiskMap; // key:域名 value:风险等级
private final JedisCluster jedisCluster;
private final String inceptorUrl = "jdbc:hive2://inceptor集群地址:10000/default";
private final String user = "xxx";
private final String pwd = "xxx";
public SensitiveDomainCacheManager(JedisCluster jedisCluster) {
this.jedisCluster = jedisCluster;
this.domainRiskMap = new HashMap<>(12_000_000);
initFullCache();
// 定时每小时全量刷新黑名单
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(this::initFullCache, 1, 1, TimeUnit.HOURS);
}
// 全量拉取Inceptor黑名单初始化缓存
private void initFullCache() {
bloomFilter = BloomFilter.create(Funnels.stringFunnel(), EXPECTED_INSERTIONS, FPP);
domainRiskMap.clear();
String sql = "SELECT domain, risk_level FROM ods_sensitive_domain_blacklist WHERE is_delete=0";
try (Connection conn = DriverManager.getConnection(inceptorUrl, user, pwd);
Statement stmt = conn.createStatement();
ResultSet rs = stmt.executeQuery(sql)) {
while (rs.next()) {
String domain = rs.getString("domain");
Integer level = rs.getInt("risk_level");
bloomFilter.put(domain);
domainRiskMap.put(domain, level);
jedisCluster.hset("sensitive:domain:map", domain, level.toString());
}
System.out.println("全量加载敏感域名完成,总数:" + domainRiskMap.size());
} catch (Exception e) {
e.printStackTrace();
}
}
// 增量更新缓存(消费Kafka增量变更时调用)
public void updateCache(String domain, Integer level, boolean isDelete) {
if (isDelete) {
domainRiskMap.remove(domain);
jedisCluster.hdel("sensitive:domain:map", domain);
} else {
bloomFilter.put(domain);
domainRiskMap.put(domain, level);
jedisCluster.hset("sensitive:domain:map", domain, level.toString());
}
}
// 匹配域名,返回风险等级,null代表不命中
public Integer matchDomain(String queryDomain) {
if (!bloomFilter.mightContain(queryDomain)) {
return null;
}
return domainRiskMap.getOrDefault(queryDomain, null);
}
}
代码3:Flink主任务(消费日志、匹配、聚合、输出告警)
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import redis.clients.jedis.JedisCluster;
import java.util.HashSet;
import java.util.Objects;
public class DnsAlarmFlinkJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(16); // 和kafka分区数对齐
env.enableCheckpointing(30000); // 30秒一次checkpoint
// 1. 构建Kafka原始DNS日志消费源
KafkaSource<String> logSource = KafkaSource.<String>builder()
.setBootstrapServers("kafka1:9092,kafka2:9092,kafka3:9092")
.setTopics("dns-log-raw")
.setGroupId("flink-dns-log-group")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<String> rawLogStream = env.fromSource(logSource, WatermarkStrategy.noWatermarks(), "dns-log-source");
// 2. 日志解析 + 敏感域名匹配
DataStream<DnsAlarm> alarmStream = rawLogStream.map(new RichMapFunction<String, DnsAlarm>() {
private SensitiveDomainCacheManager cacheManager;
@Override
public void open(org.apache.flink.configuration.Configuration parameters) throws Exception {
JedisCluster jedis = new JedisCluster(new HashSet<>()); // redis集群配置
cacheManager = new SensitiveDomainCacheManager(jedis);
}
@Override
public DnsAlarm map(String logLine) throws Exception {
// 此处替换为真实日志解析逻辑,示例模拟解析
String[] split = logLine.split("|");
DnsLog log = new DnsLog();
log.setClientIp(split[0]);
log.setQueryDomain(split[1]);
log.setDnsHost(split[2]);
log.setQueryTime(new Timestamp(Long.parseLong(split[3])));
log.setRawLog(logLine);
// 匹配敏感域名
Integer riskLevel = cacheManager.matchDomain(log.getQueryDomain());
if (riskLevel == null) return null;
// 封装告警
DnsAlarm alarm = new DnsAlarm();
alarm.setClientIp(log.getClientIp());
alarm.setSensitiveDomain(log.getQueryDomain());
alarm.setRiskLevel(riskLevel);
alarm.setDnsServerHost(log.getDnsHost());
alarm.setQueryTime(log.getQueryTime());
alarm.setRawLog(log.getRawLog());
alarm.setSyncTime(new Timestamp(System.currentTimeMillis()));
alarm.setAlarmCount(1L);
return alarm;
}
}).filter(Objects::nonNull); // 过滤未命中的日志
// 3. 5分钟滚动窗口聚合降噪,同一IP+域名合并计数
DataStream<DnsAlarm> aggAlarmStream = alarmStream
.keyBy(alarm -> alarm.getClientIp() + "_" + alarm.getSensitiveDomain())
.window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
.reduce((a1, a2) -> {
a1.setAlarmCount(a1.getAlarmCount() + a2.getAlarmCount());
return a1;
});
// 4. 告警写入Kafka dns-alarm-data
KafkaSink<DnsAlarm> alarmKafkaSink = KafkaSink.<DnsAlarm>builder()
.setBootstrapServers("kafka1:9092,kafka2:9092,kafka3:9092")
.setKafkaRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("dns-alarm-data")
.setValueSerializationSchema(new AlarmJsonSchema()) // 自定义JSON序列化
.build())
.build();
aggAlarmStream.sinkTo(alarmKafkaSink);
// 5. 独立Flink任务:消费dns-alarm-data批量写入Inceptor,这里省略独立Job代码
env.execute("DNS敏感域名实时监控任务");
}
}
代码4:Filebeat核心配置 filebeat.yml
filebeat.inputs:
– type: filestream
paths:
– /var/log/dns/*.log
close_inactive: 2m
scan_frequency: 1s
parsers:
– ndjson: # 若日志为json格式,按需切换
overwrite_keys: true
output.kafka:
hosts: ["kafka1:9092","kafka2:9092","kafka3:9092"]
topic: "dns-log-raw"
partition.round_robin:
reachable_only: false
compression: snappy
queue.mem:
events: 8192
flush.min_events: 512
flush.timeout: 1s
代码5:告警批量写入Inceptor工具类
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.Timestamp;
import java.util.List;
public class Alarm2InceptorWriter {
private static final String URL = "jdbc:hive2://inceptor集群:10000/default";
private static final String USER = "xxx";
private static final String PWD = "xxx";
private static final String INSERT_SQL = "INSERT INTO dwd_dns_query_alarm PARTITION(dt=?,hour=?) VALUES (?,?,?,?,?,?,?)";
public static void batchInsert(List<DnsAlarm> alarmList) throws Exception {
try (Connection conn = DriverManager.getConnection(URL, USER, PWD);
PreparedStatement pstmt = conn.prepareStatement(INSERT_SQL)) {
for (DnsAlarm alarm : alarmList) {
String dt = new java.text.SimpleDateFormat("yyyyMMdd").format(alarm.getQueryTime());
String hour = new java.text.SimpleDateFormat("HH").format(alarm.getQueryTime());
pstmt.setString(1, dt);
pstmt.setString(2, hour);
pstmt.setString(3, alarm.getClientIp());
pstmt.setString(4, alarm.getSensitiveDomain());
pstmt.setInt(5, alarm.getRiskLevel());
pstmt.setString(6, alarm.getDnsServerHost());
pstmt.setTimestamp(7, alarm.getQueryTime());
pstmt.setLong(8, alarm.getAlarmCount());
pstmt.setString(9, alarm.getRawLog());
pstmt.setTimestamp(10, alarm.getSyncTime());
pstmt.addBatch();
}
pstmt.executeBatch();
}
}
}
七、生产落地补充优化方案
- 域名模糊匹配拓展
若需要匹配泛域名(*.xxx.com),HashMap存储后缀,匹配时切割域名逐级后缀匹配;布隆过滤器适配后缀存储。 - 白名单机制
新增白名单IP/域名缓存,匹配敏感后再校验白名单,命中白名单直接丢弃不生成告警。 - 告警推送扩展
Flink告警流新增侧输出流,对接钉钉/短信API,高危域名实时推送告警通知。 - 冷热数据分层
Inceptor告警表30天内热数据ORC存储在线查询,超过30天迁移至HDFS冷分区,降低集群压力。 - 性能极限优化
- Flink TaskManager使用堆外内存存储布隆过滤器,减少GC停顿;
- Inceptor查询敏感黑名单时开启谓词下推、分区裁剪;
- Kafka批量发送、批量消费,减少网络IO。
- 故障应急方案
Inceptor服务不可用时,Flink自动降级使用Redis缓存的敏感域名继续匹配;Inceptor恢复后自动同步增量数据。




