欢迎光临
我们一直在努力

部署Flink、Kafka、Elasticsearch和Kibana的集群

提示:文章写完后,目录可以自动生成,如何生成可参考右边的帮助文档

目录

前言

一、前期构思

二、搭建步骤

1.搭建虚拟机集群(CentOS 7)

2.安装ES7.17.21

​编辑​编辑3.安装Flink1.19.3

4.kafka配置

5.配置Kibana

6.配置一键启动脚本

总结


前言

由于本人是网安大三牲,学校要求设计毕设,本人升本之前又是大数据专业的本科专业知识学的一知半解、于是乎决定搭建一个基于大数据流处理流量异常监测平台


一、前期构思

面对日益复杂的网络攻击与海量数据流,构建现代化的网络安全防御体系,必须满足三个核心要求:实时响应、高效处理与弹性扩展。为此,我选择了以 Apache Flink 作为实时计算引擎、Apache Kafka 作为数据采集与缓冲通道、Elasticsearch 作为存储与检索引擎,并结合 Kibana 实现可视化分析的整套技术栈。

在具体的版本选择上,原本计划采用较新的 Elasticsearch 8.13.1 与 Flink 1.19.3 组合。但在实际集成过程中发现,Flink 官方提供的 Elasticsearch 连接器(Connector)对其新版客户端的支持尚不完善,未能找到完全兼容这两个版本的稳定连接方案。为确保整套流水线的稳定与可维护性,最终决定将 Elasticsearch 版本适度回调至一个与 Flink 连接器生态兼容性更好的稳定版本。这一调整体现了在实际工程中,技术选型不仅需要考虑组件的前沿性,更要优先确保整个系统链条的可靠衔接与稳健运行。

二、搭建步骤

1.搭建虚拟机集群(CentOS 7)

由于我的电脑配置较差于是我选择配属三台虚拟机作为集群节点配置如下:

(主机1与主机2配置相同)

主机3配置如下:

为了方便后续查找与配置我选择在home目录下新建一个文件夹为bigdata(三台主机目录配相同方便后续集群分发)

2.安装ES7.17.21

  • 进入home/bigdata目录下下载ES7.17.21

    cd /home/bigdata
    wget https://artifacts.elastic.co/downloads/elasticsearch/elasticsearch-7.17.21-linux-x86_64.tar.gz

  • 解压包并重命名

    tar -zxvf elasticsearch-7.17.21-linux-x86_64.tar.gz
    mv elasticsearch-7.17.21 elasticsearch-7

    由此可见已经安装成功

  • 创建ES数据文件夹与日志文件夹,重新创建的原因也是因为防止后续ES的文件查找麻烦方便后续管理

    mkdir -p /home/bigdata/elasticsearch-7/data
    mkdir -p /home/bigdata/elasticsearch-7/logs

  • 由于 Elasticsearch 基于安全设计,禁止使用 root 超级用户直接运行,因此必须为其创建独立的专用系统账户并授权

    groupadd elastic
    useradd elastic -g elastic -p password
    chown -R elastic:elastic /home/bigdata/elasticsearch-7/

  • 修改/config/elasticsearch.yml文件 config/jvm.options

    # ======================== 集群基础配置 =========================
    cluster.name: traffic-analysis-cluster # 所有节点必须一致
    node.name: es-node-01 # 每个节点必须唯一 (es-node-02, es-node-03)

    # ======================== 路径设置 =============================
    # 前面创建的data与logs目录
    path.data: /home/bigdata/elasticsearch-7/data
    path.logs: /home/bigdata/elasticsearch-7/logs

    # ======================== 网络配置 =============================
    network.host: 0.0.0.0 # 允许外网及内网访问虽然真实环境下是不能这样的但是知识尚浅只能后续查找办法优化
    http.port: 9200 # Flink 和 Kibana 连接此端口
    transport.tcp.port: 9300 # 节点间通信端口

    # ======================== 集群发现 (ES 7 核心) =================
    # 填入三台机器的内网 IP 地址
    discovery.seed_hosts: ["192.168.x.1", "192.168.x.2", "192.168.x.3"]

    # 首次启动时参与竞选 Master 的节点名
    cluster.initial_master_nodes: ["es-node-01", "es-node-02", "es-node-03"]#前面填写的node.name

    # ======================== 性能与安全 ===========================
    # 禁用交换分区,提升性能
    bootstrap.memory_lock: true

    # 关键:暂时关闭安全校验,确保 Flink 1.19.3 能够无障碍连接
    xpack.security.enabled: false

    # 修改内存配置 (针对 JDK 8 优化)
    # 编辑 config/jvm.options,确保 -Xms 和 -Xmx 设置为你内存的一半(如 -Xms2g -Xmx2g)
    vi config/jvm.options

    需注意:elasticsearch.ym中修改项前面的#需要删除并且不能有空格不然无法生效

  • 配置结束后分发给节点2,3

    scp -r /home/bigdata/elasticsearch-7 root@Es02:/home/bigdata/
    scp -r /home/bigdata/elasticsearch-7 root@Es03:/home/bigdata/

  • 分发完成后,必须手动进入 Es02 和 Es03 的配置文件,将 node.name 改为node-02 和 node-03。如果节点名重复,集群将无法组建。

  • 切换至elastic用户进行后台启动

    su – user

    cd cd /home/bigdata/elasticsearch-7/

    ./bin/elasticsearch -d #后台运行

    3.安装Flink1.19.3

  • 单机配置下载并解压

    cd /home/bigdata/
    wget https://archive.apache.org/dist/flink/flink-1.19.3/flink-1.19.3-bin-scala_2.12.tgz
    tar -xzf flink-1.19.3-bin-scala_2.12.tgz

  • 修改 conf/config.yaml 基础配置:

    # ==============================================================================
    # JobManager 配置 (控制节点)
    # ==============================================================================
    jobmanager:
    rpc:
    address: 192.168.10.101 # 主节点 IP
    port: 6123 # 内部通信端口
    bind-host: 0.0.0.0 # 允许监听所有网卡
    memory:
    process:
    size: 1600m # JM 进程总内存
    execution:
    failover-strategy: region # 局部失败恢复策略,提高稳定性

    # ==============================================================================
    # TaskManager 配置 (计算节点)
    # ==============================================================================
    taskmanager:
    # 重要:在 101 部署时写 101;分发到 102 时请改为 102;分发到 103 时请改为 103
    host: 192.168.10.101
    bind-host: 0.0.0.0
    numberOfTaskSlots: 3 # 每个节点提供 3 个插槽 (总并发 = 3节点 * 3 = 9)
    memory:
    process:
    size: 2048m # 建议流量大时给 TM 分配 2GB 内存
    managed:
    fraction: 0.1 # 内部管理内存占比

    # ==============================================================================
    # 并行度与容错
    # ==============================================================================
    parallelism:
    default: 3 # 默认并行度,建议至少与节点数一致

    # Checkpoint 配置(确保数据不丢失)
    execution:
    checkpointing:
    interval: 30s # 每 30 秒做一次状态快照
    mode: exactly_once # 精确一次处理语义

    # ==============================================================================
    # Web UI 与 REST API
    # ==============================================================================
    rest:
    address: 192.168.10.101 # Web 访问地址
    bind-address: 0.0.0.0
    port: 8081 # 通过 http://192.168.10.101:8081 访问

  • 修改 conf/masters(配置告诉谁是主节点

    vi /home/bigdata/flink/conf/masters
    主机IP:8081

  • 修改 conf/workers(告诉都有哪些ip地址是这个集群的)

    vi /home/bigdata/flink/conf/workers
    192.168.10.101
    192.168.10.102
    192.168.10.103

  • 同ES一样分发节点2,3后修改taskmanager.host:

  • 4.kafka配置

    跟上述一样下载解压安装

    修改/conf/server.properties文件

    # ==============================================================================
    # 1. 基础配置 (每个节点必须唯一)
    # ==============================================================================
    # 节点 ID,集群内唯一。101 设为 0,102 设为 1,103 设为 2
    broker.id=0

    # ==============================================================================
    # 2. 网络与监听配置 (根据节点 IP 修改)
    # ==============================================================================
    # 监听地址
    listeners=PLAINTEXT://192.168.10.101:9092
    # 告诉客户端如何连接到当前 Broker
    advertised.listeners=PLAINTEXT://192.168.10.101:9092

    # 处理网络请求的线程数
    num.network.threads=3
    # 处理磁盘 I/O 的线程数
    num.io.threads=8

    # ==============================================================================
    # 3. 日志与存储配置
    # ==============================================================================
    # 日志文件存储路径(确保该目录存在且有写入权限)
    log.dirs=/tmp/kafka-logs

    # 每个 Topic 默认的分区数
    num.partitions=3
    # 每个分区的副本因子,建议设为 2(保证高可用)
    offsets.topic.replication.factor=2
    transaction.state.log.replication.factor=2
    transaction.state.log.min.isr=2

    # 日志保留策略(保留 168 小时 / 7 天)
    log.retention.hours=168
    # 单个日志段文件的最大大小
    log.segment.bytes=1073741824

    # ==============================================================================
    # 4. Zookeeper 集群配置 (三台机器写死一致)
    # ==============================================================================
    zookeeper.connect=192.168.10.101:2181,192.168.10.102:2181,192.168.10.103:2181
    # Zookeeper 连接超时时间
    zookeeper.connection.timeout.ms=18000

    5.配置Kibana

    下载重解压后修改config/kibana.yml文件

    # ==============================================================================
    # 基础服务配置
    # ==============================================================================
    # 默认端口 5601
    server.port: 5601

    # 允许外部浏览器访问(设为 0.0.0.0 会监听所有网卡)
    server.host: "0.0.0.0"

    # 服务的公共访问地址(用于生成分享链接)
    server.publicBaseUrl: "http://192.168.10.101:5601"

    # ==============================================================================
    # Elasticsearch 连接配置 (三节点集群)
    # ==============================================================================
    # 这里填入你的 ES 集群节点,Kibana 会自动进行负载均衡
    elasticsearch.hosts:
    – "http://192.168.10.101:9200"
    – "http://192.168.10.102:9200"
    – "http://192.168.10.103:9200"

    # 如果你在 ES 中开启了安全验证,请取消下面两行的注释并填写账号密码
    # elasticsearch.username: "kibana_system"
    # elasticsearch.password: "your_password"

    # ==============================================================================
    # 界面语言设置
    # ==============================================================================
    # 强烈建议设置为中文,方便监控项配置
    i18n.locale: "zh-CN"

    6.配置一键启动脚本

    为方便启动,使用ai配置了一个以前启动脚本

    #!/bin/bash

    # — 基础配置 —
    NODES=("192.168.10.101" "192.168.10.102" "192.168.10.103")
    KIBANA_NODE="192.168.10.101"
    # 运行 ES 的普通用户名(必须非 root)
    ES_USER="elastic"
    # ES 7 的安装路径
    ES_PATH="/home/bigdata/elasticsearch-7"
    # Kibana 的安装路径
    KIBANA_PATH="/home/bigdata/kibana-7.17"

    echo "====================================================="
    echo " 异常分析组件集群一键启动脚本 (ES 7.x 版) "
    echo "====================================================="

    # 1. ZooKeeper & 2. Kafka (保持你原有的逻辑)
    echo ">>> [1/5] 启动 ZooKeeper & Kafka 集群…"
    for HOST in "${NODES[@]}"; do
    ssh $HOST "source /etc/profile; /home/bigdata/kafka/bin/zookeeper-server-start.sh -daemon /home/bigdata/kafka/config/zookeeper.properties"
    ssh $HOST "source /etc/profile; /home/bigdata/kafka/bin/kafka-server-start.sh -daemon /home/bigdata/kafka/config/server.properties"
    echo " – $HOST: ZK & Kafka 已启动"
    done

    # 3. Elasticsearch (降级后的核心改动)
    echo ">>> [3/5] 启动 Elasticsearch 7 集群…"
    for HOST in "${NODES[@]}"; do
    # 使用 su 切换到普通用户执行启动命令,-d 为后台运行
    ssh $HOST "su – $ES_USER -c 'nohup $ES_PATH/bin/elasticsearch -d > /dev/null 2>&1 &'"
    echo " – $HOST: ES 指令已以 $ES_USER 身份下发"
    done

    echo "等待 ES 集群初始化 (40s)…"
    sleep 40
    echo "====================================================="
    echo " 异常分析组件集群一键启动脚本 (ES 7.x 版) "
    echo "====================================================="

    # 1. ZooKeeper & 2. Kafka
    echo ">>> [1/5] 启动 ZooKeeper & Kafka 集群…"
    for HOST in "${NODES[@]}"; do
    ssh $HOST "source /etc/profile; /home/bigdata/kafka/bin/zookeeper-server-start.sh -daemon /home/bigdata/kafka/config/zookeeper.properties"
    ssh $HOST "source /etc/profile; /home/bigdata/kafka/bin/kafka-server-start.sh -daemon /home/bigdata/kafka/config/server.properties"
    echo " – $HOST: ZK & Kafka 已启动"
    done

    # 3. Elasticsearch (降级后的核心改动)
    echo ">>> [3/5] 启动 Elasticsearch 7 集群…"
    for HOST in "${NODES[@]}"; do
    # 使用 su 切换到普通用户执行启动命令,-d 为后台运行
    ssh $HOST "su – $ES_USER -c 'nohup $ES_PATH/bin/elasticsearch -d > /dev/null 2>&1 &'"
    echo " – $HOST: ES 指令已以 $ES_USER 身份下发"

    赋予权限  chmod +x cluster-start.sh 

    启动后检查服务

    登录192.168.10.101:8081,192.168.10.101:5601

    192.168.10.101:8081

    192.168.10.101:5601

    总结

    至此部署成功,一个大致框架构成,后续想的是是否可以加入机器学习监测,目前仅使用过py脚本生成相对格式的流量访问但是没有真实数据,下一步打算使用权威发布的数据集进行实验查看是否可行,但是更希望的是可以采集真实流量。其实也不确定可行性高不高

    赞(0)
    未经允许不得转载:171主机测评 » 部署Flink、Kafka、Elasticsearch和Kibana的集群
    分享到: 更多 (0)

    评论 抢沙发

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