Apache Flume 教程:什么是 Archi架构与 Hadoop 示例

⚡ 智能摘要

Apache Flume 是一种分布式服务,用于收集、聚合和将大量日志数据移动到 HDFS 中,它围绕着将源、通道和接收器链接在一起的代理构建。

  • 🔘 特工解剖结构: 每个 Flume 代理都是一个 JVM 进程,包含一个源、一个或多个通道和一个接收器。
  • ☑️ 可靠性: 尽力而为的交付方式不允许任何节点发生故障;端到端交付方式可以承受多个节点发生故障。
  • 建立: 自定义源类编译成 JAR 文件,然后将其放入 Flume lib 目录中。
  • 🧪 配置: 一个属性文件指定源、通道和目标,并设置 HDFS 路径和​​滚动限制。
  • 🛠️ 发射: 使用 flume-ng 代理启动管道,命名代理并指向 flume.conf。
  • ⚠️ 过时示例: Twitter v1.1 流媒体端点已于 2023 年 3 月关闭,因此请将此练习视为自定义源模式。

Apache Flume 教程,涵盖代理架构、配置和 Hadoop 流式传输示例

Hadoop 中的 Apache Flume 是什么?

Apache Flume 是一个可靠的分布式系统,用于收集、聚合和传输海量日志数据。它拥有基于流式数据流的简洁而灵活的架构。Apache Flume 用于从 Web 服务器收集日志文件中的日志数据,并将其聚合为 高密度文件系统 进行分析。

Flume 在 Hadoop 中支持多种数据源,包括:

  • 'tail'(它将本地文件中的数据通过管道传输到 HDFS,类似于 Unix 命令 'tail')
  • 系统日志
  • Apache log4j (这使得 Java 应用程序通过 Flume 将事件写入 HDFS 中的文件)。

当前版本是 Flume 1.11.0于2022年10月25日发布,可从以下渠道获取: Apache Flume 下载页面本教程是基于 1.4.0 版本编写的,因此以下几个步骤中会注明当前版本的行为有所不同。

水槽 Archi质地

Flume代理是一种 JVM 该流程包含三个组件——Flume 源、Flume 通道和 Flume 汇——事件在外部源发起后,会依次通过这三个组件进行传播。下图展示了它们的连接方式。

Flume架构图展示了一个代理,该代理包含源、通道和接收器,并为HDFS提供数据。

  1. Flume 源会接收来自外部源(Web 服务器)的事件。外部源以目标源可识别的格式向 Flume 源发送事件。
  2. Flume 源接收事件并将其存储到一个或多个通道中。通道充当存储库,保存事件直到被 Flume 接收器消费。该通道可以使用本地文件系统来存储这些事件。
  3. Flume 接收器会将事件从通道中移除并存储到外部存储库(例如 HDFS)。如果存在多个 Flume 代理,Flume 接收器会将事件转发给流程中下一个代理的 Flume 源。

水槽的一些重要特征

  • Flume 采用基于流式数据流的灵活设计,具有容错性和鲁棒性,并提供多种故障转移和恢复机制。Flume 提供不同级别的可靠性,包括: “尽力递送”“端到端交付”. 尽力而为 不能容忍任何 Flume 节点故障,而 端到端交付 即使多个节点发生故障,也能保证交付。
  • Flume负责在数据源和目标数据之间传输数据。这种数据采集可以是定时的,也可以是事件驱动的。Flume拥有自己的查询处理引擎,可以轻松地对每批新数据进行转换,然后再将其发送到目标数据。
  • 可能存在 水槽 包括HDFS和 HBase的Flume还可以传输事件数据,例如网络流量数据、社交媒体网站生成的数据和电子邮件消息。

Flume、库和源代码设置

在开始实际操作之前,请确保已安装 Hadoop;如果没有,请按照以下步骤进行操作。 如何安装Hadoop 首先,将用户更改为“hduser”(配置 Hadoop 时使用的 ID;您可以切换到您自己的 Hadoop 配置期间使用的用户 ID)。

在 Flume 安装开始之前,终端会将 Linux 用户切换到 hduser。

步骤1) 创建一个名为“FlumeTutorial”的新目录。

sudo mkdir FlumeTutorial
  1. 授予读取、写入和执行权限。
    sudo chmod -R 777 FlumeTutorial
  2. 复制文件 我的TwitterSource.javaMyTwitterSourceForFlume.java 进入此目录。

从这里下载输入文件

检查以下所有文件的文件权限,如果缺少“读取”权限,请授予该权限。

终端列出已下载文件的权限 Java 源文件

步骤2) 从此处下载“Apache Flume” https://flume.apache.org/download.html.

本 Flume 教程中使用了 Apache Flume 1.4.0。

Apache Flume 下载页面显示二进制 tarball 链接以供选择

接下来,点击进入镜像页面。

点击 Flume tarball 链接后到达的 Apache 镜像页面

步骤3) 将下载的 tar 包复制到您选择的目录中并执行trac使用以下命令读取内容。

sudo tar -xvf apache-flume-1.4.0-bin.tar.gz

终端出口trac使用 sudo tar -xvf 命令解压 Flume tarball 文件

这将创建一个名为 apache-flume-1.4.0-bin 的新目录,并执行trac将文件放入其中。该目录被称为文章其余部分。

步骤4) Flume库配置。将twitter4j-core-4.0.1.jar、flume-ng-configuration-1.4.0.jar、flume-ng-core-1.4.0.jar和flume-ng-sdk-1.4.0.jar复制到

/lib/

复制的 JAR 文件可能全部或部分具有执行权限,这会导致代码编译出现问题,因此需要撤销该权限。就我而言,twitter4j-core-4.0.1.jar 具有执行权限。我按如下方式撤销了该权限。

sudo chmod -x twitter4j-core-4.0.1.jar

终端正在撤销对 twitter4j 核心 JAR 文件的执行权限

之后,以下命令将对 twitter4j-core-4.0.1.jar 授予所有人“读取”权限。

sudo chmod +rrr /usr/local/apache-flume-1.4.0-bin/lib/twitter4j-core-4.0.1.jar

请注意,我从以下位置下载了 twitter4j-core-4.0.1.jar: Maven 仓库以及所有 Flume JAR 文件,例如 flume-ng-*-1.4.0.jar,来自 org.apache.flume 制品.

使用 Flume 从 Twitter 加载数据

步骤1) 进入包含源代码文件的目录。

步骤2) 设置 CLASSPATH 以包含/lib/* 和 ~/FlumeTutorial/flume/mytwittersource/*。

export CLASSPATH="/usr/local/apache-flume-1.4.0-bin/lib/*:~/FlumeTutorial/flume/mytwittersource/*"

终端导出指向 Flume 库和源代码目录的 CLASSPATH

步骤3) 使用以下命令编译源代码。

javac -d . MyTwitterSourceForFlume.java MyTwitterSource.java

终端编译这两个 Java 使用 javac 生成源文件

步骤4) 创建 JAR 文件。首先,使用您选择的文本编辑器创建一个名为 Manifest.txt 的文件,并将以下行添加到该文件中。

Main-Class: flume.mytwittersource.MyTwitterSourceForFlume

这里 flume.mytwittersource.MyTwitterSourceForFlume 是主类的名称。请注意,您需要在本行末尾按回车键,如下所示。

用文本编辑器打开 Manifest.txt 文件,并找到 Main-Class 条目。

现在,按如下方式创建 JAR 文件“MyTwitterSourceForFlume.jar”。

jar cfm MyTwitterSourceForFlume.jar Manifest.txt flume/mytwittersource/*.class

终端将编译后的类打包到 MyTwitterSourceForFlume.jar 中。

步骤5) 将此 JAR 文件复制到/lib/。

sudo cp MyTwitterSourceForFlume.jar <Flume Installation Directory>/lib/

终端正在将自定义源 JAR 文件复制到 Flume lib 目录中

步骤6) 进入 Flume 的配置目录, /conf。

如果 flume.conf 不存在,请复制 flume-conf.properties.template 并将其重命名为 flume.conf。

sudo cp flume-conf.properties.template flume.conf

终端正在将 flume-conf.properties.template 复制到 flume.conf

如果 flume-env.sh 不存在,请复制 flume-env.sh.template 并将其重命名为 flume-env.sh。

sudo cp flume-env.sh.template flume-env.sh

终端正在将 flume-env.sh.template 复制到 flume-env.sh

创建 Twitter 应用程序

请先阅读此文。 v1.1 流媒体 statuses/filter twitter4j 4.0.1 所需的端点已于 2023 年 3 月 9 日停用,取而代之的是 API v2 过滤流,现在位于 developer.x.com此功能位于付费层级之后。请将以下屏幕视为自定义源模式,然后将同一代理指向文件、可执行文件或 Kafka 源。

步骤1) 登录开发者门户,创建 Twitter 应用。

用于访问应用程序列表的 Twitter 开发者登录页面

登录后显示的是 Twitter 开发者账号主页

步骤2) 前往“我的应用程序”(点击右上角的“蛋”按钮后,此选项将下拉显示)。

Twitter开发者门户的“我的应用”页面

步骤3) 点击“创建新应用”来创建一个新应用。

步骤4) 请填写申请详情,包括申请名称、申请描述和网站地址。您可以参考每个输入框下方的说明。

Twitter 应用创建表单,包含名称、描述和网站字段

步骤5) 向下滚动页面,勾选“是,我同意”接受条款,然后点击“创建您的 Twitter 应用程序”按钮。

Twitter 表单底部的“条款”复选框和“创建申请”按钮

步骤6) 在新建的应用程序窗口中,转到“API 密钥”选项卡,向下滚动页面,然后单击“创建我的访问令牌”按钮。

在访问令牌存在之前,新版 Twitter 应用程序的“API 密钥”选项卡

使用“创建我的访问令牌”按钮后,将显示访问令牌详细信息。

步骤7) 刷新页面。

步骤8) 点击“测试 OAuth”。这将显示应用程序的“OAuth”设置。

测试 OAuth 屏幕,显示应用程序的 OAuth 设置。

步骤9) 请使用以下 OAuth 设置修改 'flume.conf' 文件。修改 'flume.conf' 文件的步骤如下。

OAuth 设置中列出了消费者密钥、消费者密钥和访问令牌的值

我们需要复制消费者密钥、消费者密钥、访问令牌和访问令牌密钥,以便更新“flume.conf”。

注意:这些值属于用户所有,因此属于机密信息,不应共享。

修改 flume.conf 文件

步骤1) 以写入模式打开“flume.conf”文件,并设置以下参数的值。

sudo gedit flume.conf

复制以下内容。

MyTwitAgent.sources = Twitter
MyTwitAgent.channels = MemChannel
MyTwitAgent.sinks = HDFS
MyTwitAgent.sources.Twitter.type = flume.mytwittersource.MyTwitterSourceForFlume
MyTwitAgent.sources.Twitter.channels = MemChannel
MyTwitAgent.sources.Twitter.consumerKey = <Copy consumer key value from Twitter App>
MyTwitAgent.sources.Twitter.consumerSecret = <Copy consumer secret value from Twitter App>
MyTwitAgent.sources.Twitter.accessToken = <Copy access token value from Twitter App>
MyTwitAgent.sources.Twitter.accessTokenSecret = <Copy access token secret value from Twitter App>
MyTwitAgent.sources.Twitter.keywords = guru99
MyTwitAgent.sinks.HDFS.channel = MemChannel
MyTwitAgent.sinks.HDFS.type = hdfs
MyTwitAgent.sinks.HDFS.hdfs.path = hdfs://localhost:54310/user/hduser/flume/tweets/
MyTwitAgent.sinks.HDFS.hdfs.fileType = DataStream
MyTwitAgent.sinks.HDFS.hdfs.writeFormat = Text
MyTwitAgent.sinks.HDFS.hdfs.batchSize = 1000
MyTwitAgent.sinks.HDFS.hdfs.rollSize = 0
MyTwitAgent.sinks.HDFS.hdfs.rollCount = 10000
MyTwitAgent.channels.MemChannel.type = memory
MyTwitAgent.channels.MemChannel.capacity = 10000
MyTwitAgent.channels.MemChannel.transactionCapacity = 1000

使用编辑器打开 flume.conf 文件,其中包含 MyTwitAgent 的源、通道和接收器属性。

步骤2) 另外,请按如下方式设置 TwitterAgent.sinks.HDFS.hdfs.path。

TwitterAgent.sinks.HDFS.hdfs.path = hdfs:// : / /flume/推文/

HDFS 接收器 hdfs.path 属性设置为主机名、端口号和 HDFS 主目录

查找,和请查看 $HADOOP_HOME/etc/hadoop/core-site.xml 中设置的参数“fs.defaultFS”的值,如下所示。

core-site.xml 文件中的 fs.defaultFS 属性提供了主机名和端口。

步骤3) 为了在数据到达时将其刷新到 HDFS,如果存在以下条目,请将其删除。

TwitterAgent.sinks.HDFS.hdfs.rollInterval = 600

示例:使用 Flume 传输 Twitter 数据

步骤1) 以写入模式打开“flume-env.sh”文件,并设置以下参数的值。

JAVA_HOME=<Installation directory of Java>
FLUME_CLASSPATH="<Flume Installation Directory>/lib/MyTwitterSourceForFlume.jar"

使用已设置 JAVA_HOME 和 FLUME_CLASSPATH 的编辑器打开 flume-env.sh 文件。

步骤2) 启动Hadoop。

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

步骤3) Flume tar 包中的两个 JAR 文件与 Hadoop 2.2.0 不兼容,因此在本 Apache Flume 示例中,我们将按照以下步骤使 Flume 与 Hadoop 2.2.0 兼容。此 JAR 文件替换是 1.4.0 版本时代的修复;Flume 1.11.0 已经包含了最新的 protobuf 和 Guava 构建版本,因此现代 tar 包通常不需要此操作。

a. 将 protobuf-java-2.4.1.jar 从“请先进入/lib目录。

光盘/lib

sudo mv protobuf-java-2.4.1.jar ~/

终端正在将 protobuf-java-2.4.1.jar 从 Flume lib 目录中移出。

b. 找到如下所示的 JAR 文件“guava”。

find . -name "guava*"

使用终端 find 命令查找捆绑的 Guava JAR 文件

将 guava-10.0.1.jar 从 ' 中移出。 /lib'。

sudo mv guava-10.0.1.jar ~/

终端正在将 guava-10.0.1.jar 从 Flume lib 目录中移出。

c. 从以下位置下载 guava-17.0.jar Maven 仓库,如下所示。

Guava 17.0 的 Maven 仓库页面,可下载替换 JAR 文件。

现在,将下载的 JAR 文件复制到“ /lib'。

步骤4) 前往“ /bin' 并按如下方式启动 Flume。

./flume-ng agent -n MyTwitAgent -c conf -f <Flume Installation Directory>/conf/flume.conf

使用 flume-ng 命令在终端启动名为 MyTwitAgent 的 Flume 代理。

Flume 获取推文的命令提示符窗口如下所示。

命令提示符显示 Flume 代理正在抓取推文并将其写入 HDFS。

从命令窗口的消息中我们可以看到,输出被写入了 /user/hduser/flume/tweets/ 目录。现在,使用网页浏览器打开这个目录。

步骤5) 要查看数据加载结果,请在浏览器中打开 http://localhost:50070/,浏览文件系统,然后转到数据已加载的目录,即

/flume/推文/

Hadoop 2 的 NameNode Web UI 使用端口 50070;Hadoop 3 将同一页面移至端口 9870。

HDFS浏览器显示了flume/tweets目录及其加载的推文文件

水槽是摄入量的一半: 勺子 Flume批量导入表格,然后传输事件流。 or 蜂房 对文件进行格式化和 乌兹 安排链条。另请参阅 大数据分析工具, MapReduce 连接和计数器拓蓝.

常见问题

并非如描述的那样。v1.1 流式状态/筛选接口已于 2023 年 3 月 9 日停用,替代的 API v2 需要付费才能使用。Flume 的机制仍然适用于自定义源代码。

模型首先建立基线正态对数体积和消息形状,然后标记固定阈值无法识别的偏差。它们还会对重复堆栈进行聚类。 trac将多个事件合并为一个单一事件,并拟定可能原因,从而缩短分诊时间。

Copilot 可以快速生成源、通道和接收器模块,但它会自行命名属性并混用不同版本的 Flume。启动代理之前,请务必对照您 Flume 版本的用户指南检查每个键值。

Flume 将日志等事件数据持续不断地流式传输到 HDFS。Sqoop 则按计划批量地在关系数据库和 Hadoop 之间迁移结构化表。它们分别负责数据摄取的不同环节,配合使用效果极佳。

内存通道速度最快,但如果代理程序崩溃,则会丢失缓存的事件。文件通道将数据写入磁盘,重启后数据仍然保留,但吞吐量较低。对于无法重新发送的数据,建议优先选择持久化方案。

Kafka 现在通常是默认选择,因为它能保留数据并服务于众多消费者。Flume 1.11.0(2022 年 10 月发布)仍然适用于将简单的单向日志收集到 HDFS 中。

几乎总是 JAR 包冲突:Flume 压缩包捆绑了它自己的 Guava 和 protobuf 版本,这与 Hadoop 加载的版本冲突。移除旧的捆绑 JAR 包通常可以解决问题。

它们决定接收器何时关闭文件并打开新文件:rollSize 按字节数计算,rollCount 按事件计数计算,rollInterval 按秒数计算。零表示禁用该特定触发器。

总结一下这篇文章: