欢迎光临
我们一直在努力

实时监控DNS敏感域名全方案分析与落地实现

实时监控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(核心)

核心能力:

  • 消费Kafka原始DNS日志流,解析客户端IP、查询域名、服务器标识、时间;
  • 增量同步Inceptor千万级敏感域名库,构建内存布隆过滤器+HashMap双层缓存;
  • 实时域名匹配,命中则生成告警;
  • 窗口聚合降噪(5分钟滚动窗口,同一IP+域名合并一条告警);
  • 告警写入Kafka告警Topic,同步落地Inceptor告警表;
  • 状态持久化(Checkpoint),任务重启不丢失数据。
    • 为什么不用Spark Streaming:Spark微批分钟级延迟,无法满足实时秒级监控;Flink原生流处理,低延迟。
    4. 缓存层:Redis(配套Flink)
    • 用途1:存储全量敏感域名HashSet,Flink本地缓存失效兜底;
    • 用途2:告警降噪去重,窗口内重复命中IP+域名计数;
    • 用途3:缓存最新一批敏感域名,加速增量同步;
    5. 存储层
  • 星环Inceptor
    • 两张核心表:敏感域名黑名单表、DNS告警明细表;
    • 优势:兼容Hive SQL,海量结构化数据离线统计、多维度查询;千万级黑名单可增量抽取;
  • HDFS:原始DNS日志冷归档,按需回溯原始日志;
  • Redis:实时告警快速查询、降噪缓存;
  • 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精确匹配:布隆命中后校验,消除误判;
  • 缓存兜底:同步更新Redis全局敏感域名集合,Flink任务重启时快速加载,无需重新拉取千万级全量。
  • 禁止方案:每条DNS日志直接JDBC查询Inceptor,QPS上万时数据库直接打垮,延迟分钟级,不满足实时。

    2. 完整数据流处理步骤

  • Filebeat监听DNS日志文件,实时读取新增行,附加DNS服务器标识、采集时间;过滤无效日志,发送Kafka dns-log-raw;
  • Flink消费原始日志流:
    • 步骤1:日志清洗解析,提取client_ip、query_domain、dns_host、log_time,过滤非法空域名;
    • 步骤2:读取本地内存敏感域名缓存,执行布隆+Hash双层匹配;
    • 步骤3:命中敏感域名 → 封装告警实体(IP、域名、服务器、时间、日志原文、敏感域名等级);
    • 步骤4:5分钟滚动窗口聚合,同一IP+同一域名合并计数,实现告警降噪,避免刷屏;
    • 步骤5:聚合后告警写入Kafka dns-alarm-data;
  • Flink另一分支消费告警Topic,批量写入Inceptor告警明细表(批量Insert提升写入性能);
  • 原始日志同步归档HDFS,保存30天用于溯源;
  • Prometheus采集全链路指标,Grafana可视化监控延迟、吞吐量。
  • 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恢复后自动同步增量数据。
    赞(0)
    未经允许不得转载:171主机测评 » 实时监控DNS敏感域名全方案分析与落地实现
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址