ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flume到Kafka日志采集链路配置详解与踩坑实录

Flume到Kafka日志采集链路配置详解与踩坑实录 说实话我学尚硅谷这套数仓实战课的时候前几节搭 Hadoop、配 Zookeeper、把 Kafka 集群跑起来感觉都还挺顺的直到第 8 节——配置 Flume 把日志文件丢进 Kafka才对“生产环境下的日志采集到底怎么玩”有了实感。这一节是整个实时数仓任务的数据入口业务系统产生日志文件Flume 趴在文件系统上盯着日志增长读出来之后通过 Kafka 这个中间缓冲层交给下游的 Flink 或其他消费端。如果你正在学数仓搭建或者想在自己电脑上把离线/实时数仓链路跑通这一节是绕不过去的核心环节。这篇文章会用我当时的配置文件和踩坑记录把从写配置到跑通全过程的细节掰开讲清楚也包括那几个让我卡了一晚上的坑。1. 数仓日志链路里Flume 到 Kafka 这一环是干嘛的1.1 课程进度到这儿整个采集链路长什么样这一节在整个数仓项目里的位置简单说就是“实时数仓的快递第一公里”。整个链路大概是这样的业务服务器上的应用会不断往磁盘写日志文件比如 SpringBoot 的 logback 输出、Nginx 的 access.log。这些日志文件分散在多台机器上不会自己跑到你的数仓里去。Flume 做的事情就是站在那里盯着日志目录文件一旦有新内容就增量读出来包装成 Flume 内部的 Event 对象交给 Channel 缓冲最后由 Kafka Sink 批量写到 Kafka 的某个 topic 里。Kafka 拿到这批数据之后下游就有两条路可以走。一条是实时的用 Flink 或者 SparkStreaming 去消费 Kafka 里的数据做清洗、ETL、实时指标计算另一条是离线的定时批处理把 Kafka 里的数据落到 HDFS 或 Hive 表再进离线数仓。所以说Flume 到 Kafka 这一段是整个日志数据的“交通枢纽”之前学的 Kafka 在这里第一次真正接到数据后面学的 Flink 也从这里开始找数据。1.2 为什么日志不直接发 Kafka非要绕一道 Flume我第一次看到这个架构时也有疑问业务系统里直接调 KafkaProducer把日志发到 Kafka 不就完事了吗何苦再多塞一个 Flume后来想明白了几点。第一业务系统的日志通常是以文件形式存在的而且会按天切分、按大小切分。采集日志需要监控文件变化、按行读取、断点续传这些逻辑如果塞给业务代码会让业务系统变得很重。Flume 就是专门干这个的它把“日志怎么被读出来”这件事从业务里剥离开。业务系统只负责写文件剩下的采集和投递都交给 Flume两边解耦业务系统不会因为 Kafka 抖动就一起卡住。第二Kafka 是消息的高速公路但得有人把货物从仓库搬上高速。直接让每个应用服务直连 Kafka一旦 Kafka 集群出问题producer 的重试机制会把业务线程堵住甚至会拖垮应用。Flume 作为独立进程即使 Kafka 故障日志也只会堆在本地或者 Flume 的 Channel 里恢复以后还能继续补送这个缓冲隔离作用在真实生产里非常关键。第三Flume 支持多级拓扑。你可以让采集层的 Flume 先把日志汇总到一台汇聚节点再由汇聚节点的 Flume 做多路分发一份数据同时进 Kafka 和 HDFS。这个能力直接复用后面项目里“一份数据多次消费”的需求非常实用。1.3 学这节之前先确认这些环境是好的动手之前别急着写配置先把底下的环境确认一遍。我当时以为 Kafka 装好了就能直接用结果配置写完之后半天查不出问题最后才发现是 Kafka 客户端版本不兼容。这一节的硬前提至少有这几个。Kafka 集群要能正常创建 topic、能生产消费。验证命令很简单kafka-topics.sh --list --bootstrap-server hadoop102:9092能列出当前所有 topic 列表说明 broker 正常。再拿 console-consumer 测一下能不能消费基本就确认 Kafka 可用。Zookeeper 要正常。虽然 Kafka 新版本可以配 KRaft 模式但课程环境用的还是 Zookeeper 做协调Flume 本身不直接连 ZK但 Kafka broker 的状态依赖它所以 ZK 挂了 Kafka 整个链路都会异常。Flume 已经解压好并且版本注意一下。我用的 Flume 是 1.9.0它在发布那会儿自带的 Kafka 客户端 jar 包版本偏旧而课程环境里的 Kafka 是 2.x两者放一块就容易出那个经典的 ClassNotFoundException。这个问题后面单独开一节讲因为太典型了几乎每个人都会碰到。主机名解析也要确认。配置里写的是 hadoop102、hadoop103、hadoop104 这种主机名不是 IP所以 /etc/hosts 里必须配好映射。如果你自己练习用的是单机那写 localhost 或者 192.168.x.x 都行但一定要保证 Flume 和 Kafka 之间网络能通。2. 配置 Flume 之前先吃透 Source、Channel、Sink 这三个组件2.1 Taildir Source 是主角positionFile 机制要弄明白Flume 的 Source 负责“从哪读数据”。日志采集场景下最常用的是 Taildir Source它专门用来监控一批文件按行读取新追加的内容。为什么不用 Exec SourceExec Source 本质上是启动一个外部命令比如tail -F然后读取命令输出。这个方案的问题在于如果 Flume 进程重启tail 命令的历史位置就丢了可能重复读或者漏读。SpoolDir Source 也不合适它监控整个目录文件一旦放入就会被读走然后改名适合“一次性放入文件”的场景不适合持续追加的日志文件。Taildir Source 的核心在于 positionFile 机制。它会记录每个被监控文件当前读到哪个字节位置这些信息存在一个 JSON 文件里。Flume 每次启动时会读取这个文件从上次记录的位置继续往下读而不是从头再来。这个能力对日志采集来说太重要了因为日志文件是持续追加的Flume 进程也可能因为各种原因重启如果每次重启都重复读取历史数据下游的 Kafka 里就会出现大量脏数据。positionFile 的位置默认是~/.flume/taildir_position.json。这个默认路径有个小坑如果 Flume 的启动用户不一样或者 HOME 目录有变化positionFile 就换了地方Flume 会重新扫描文件导致重复消费。我建议显式指定到一个固定目录比如后面配置里写的/opt/module/flume/data/taildir_position.json并且确保这个目录提前创建好。2.2 Channel 选型内存快但会丢文件稳但慢Channel 是 Flume 内部的数据缓冲区连接 Source 和 Sink。它决定了数据在 Flume 内部的可靠性等级。常用的就两种Memory Channel 和 File Channel。Memory Channel 就是把 Event 存在内存里读写速度极快吞吐高。缺点也很明显进程一旦崩溃Channel 里还没发出去的数据就全丢了。File Channel 则先把 Event 写入本地磁盘Flume 崩溃后可以按日志恢复数据可靠性高但磁盘 IO 会成为瓶颈吞吐比内存模式低不少。这节课的日志采集场景选 Memory Channel 完全够用。因为日志数据本身允许少量丢失而且 Flume 到 Kafka 这条链路下游 Kafka 也有自己的持久化机制没必要让 Flume 再做一层磁盘缓冲否则白白牺牲性能。Memory Channel 有几个参数值得认真配。capacity是 Channel 里最多能存多少个 Event默认值只有 100这个值非常小Kafka 或下游稍微卡顿一下Channel 就满了满了之后 Source 就会停止读取日志文件造成反压。我一般都把它调到 10000。transactionCapacity是每次事务能处理的最大 Event 数量必须小于 capacity通常取 capacity 的 5%~10%比如 10000 配 500。这两个参数不调对你会在演示阶段看到 Flume 日志里频繁出现 Channel full 的报错但你可能根本意识不到是这里的问题。2.3 Kafka Sink 的配置关键点Sink 决定数据往哪走。这节课的主角是 Kafka Sink它的完整类名是org.apache.flume.sink.kafka.KafkaSink。底层就是包装了一个 KafkaProducer从 Channel 里批量拿 Event然后发送到指定的 Kafka topic。配置 Kafka Sink 时最核心的是三个东西。第一是kafka.bootstrap.servers告诉 Sink 怎么找到 Kafka 集群。一般至少写两个 broker 地址保证某个节点故障时还能连上。第二是kafka.topic指定目标主题。第三是kafka.flumeBatchSize表示 Sink 从 Channel 里取多少个 Event 再批量发给 Kafka。这个值配得太大会导致一批数据在本地攒很久延迟上升配得太小KafkaProducer 会频繁发送小批次数据吞吐下降。我个人实践下来100 到 500 之间是一个比较平衡的范围。还有一个 offset 相关的点Kafka Sink 发送的是批量事件它会在 KafkaProducer 发送成功后才把对应的事务提交给 Channel。所以 Channel 里的数据不会因为 Kafka 短暂不可用而立刻丢失但会造成积压。如果你看到 Channel 的 put 计数远大于 take 计数基本可以判断是 Sink 写不进 Kafka优先查 Kafka 侧的问题。另外Kafka Sink 有个隐藏的高级特性如果某个 Event 的 header 里带了topic字段KafkaSink 会优先用 header 里的 topic而忽略配置里的固定 topic。这个特性可以用来实现“一份 Flume 配置把不同类型日志写到不同 Kafka topic”但大多数场景下我们只是用配置里的固定 topic所以了解即可。3. 实操从零写一份可用的 Flume 配置并跑通3.1 准备工作创建主题、准备日志目录、建好 positionFile 位置在写配置文件之前先把基础设施准备好。第一步在 Kafka 里创建 topic我用的课程环境是三台节点副本数设为 2分区数设为 3对应命令kafka-topics.sh --create \ --bootstrap-server hadoop102:9092 \ --replication-factor 2 \ --partitions 3 \ --topic topic_log如果是单机练习replication-factor 只能设 1否则副本分配不上去命令行会报错。这个细节容易忽略但报错信息会直接告诉你副本数大于 broker 数属于正常环境限制。第二步创建日志目录和 Flume 的 positionFile 目录mkdir -p /opt/module/applog/log mkdir -p /opt/module/flume/data touch /opt/module/applog/log/app.2024-01-01.log这里我建议一开始就 touch 出一个带具体日期命名的日志文件因为后面的 filegroups 配置要匹配这种命名规则。你用app.log命名也行关键是正则要能对上。3.2 核心配置文件逐行拆解配置文件放在/opt/module/flume/conf/flume_kafka.conf完整的配置长这样a1.sources r1 a1.channels c1 a1.sinks k1 a1.sources.r1.type TAILDIR a1.sources.r1.positionFile /opt/module/flume/data/taildir_position.json a1.sources.r1.filegroups f1 a1.sources.r1.filegroups.f1 /opt/module/applog/log/app.*.log a1.sources.r1.fileHeader true a1.channels.c1.type memory a1.channels.c1.capacity 10000 a1.channels.c1.transactionCapacity 500 a1.sinks.k1.type org.apache.flume.sink.kafka.KafkaSink a1.sinks.k1.kafka.bootstrap.servers hadoop102:9092,hadoop103:9092,hadoop104:9092 a1.sinks.k1.kafka.topic topic_log a1.sinks.k1.kafka.flumeBatchSize 100 a1.sinks.k1.kafka.producer.acks 1 a1.sources.r1.channels c1 a1.sinks.k1.channel c1逐段解释一下为什么这么写。第一段是给三个组件起名字。我习惯用r1/c1/k1这种命名简洁不容易看花眼。注意组件名字必须和后面配置里的引用完全一致少写一个s或者大小写错了Flume 启动就会直接报配置语法错误。第二段配置 Taildir Source。filegroups是文件组你可以把具有相同模式的多组文件放在不同组里。filegroups.f1就是具体匹配规则了app.*.log这个正则会匹配app.log、app.2024-01-01.log等所有以 app 开头、以 .log 结尾的文件。fileHeader true表示给读到的每个 Event 的 header 里加上当前文件名这个字段在后续做数据来源识别时很有用。第三段是 Memory Channel 配置容量我调到了 10000。这个值不是拍脑袋定的我前面说过的 reasoning 就是防反压。transactionCapacity 设 500也就是每个事务最多处理 500 个 Event比 capacity 小一个数量级保证事务不会一次性把 Channel 塞满。第四段是 Kafka Sink。acks 1表示 Kafka 的 leader 收到数据就认为发送成功。日志场景下这个级别够了。如果你要求更严格的可靠性可以设成-1意思是所有副本都写入才算成功但吞吐会打折。学习阶段不用纠结这个参数知道它是控制可靠性和性能权衡的开关就行。最后两行是“接线”的固定句式把 Source 和 Sink 分别连接到同一个 Channel 上千万不能漏掉。很多新手报错说数据读进来了但发不出去最后发现 Source 接了 ChannelSink 没接数据在 Channel 里堆着Flume 日志里 Sink 的计数一直是 0。3.3 启动 Flume 并验证 Kafka 是否收到日志配置文件就绪后先别急着后台启动我推荐前台跑一次把日志打到控制台确认没有报错后再转后台。前台启动命令/opt/module/flume/bin/flume-ng agent \ --name a1 \ --conf /opt/module/flume/conf \ --conf-file /opt/module/flume/conf/flume_kafka.conf \ -Dflume.root.loggerINFO,console--name a1必须和配置文件里的a1一致否则 Flume 会找不到 agent 配置。-Dflume.root.loggerINFO,console让日志输出到控制台方便实时观察。启动之后打开另一个终端给日志文件追加一条模拟数据echo {event:page_view,user_id:10001,time:2024-01-01 10:30:00} /opt/module/applog/log/app.2024-01-01.log然后在第三个终端启动 Kafka 消费者看看能不能收到这条数据kafka-console-consumer.sh \ --bootstrap-server hadoop102:9092 \ --from-beginning \ --topic topic_log如果消费端打印出了刚才那条 JSON说明全链路已经通了。这一步我第一次跑的时候等了大概两三秒才有数据因为 Flume 对文件变化的轮询不是实时的加上 Sink 批量发送也有延迟稍微等一下是正常的。还有一个验证细节是看 Flume 控制台日志里的计数器。你会在日志里看到类似Sink: k1 ... EventPutSuccessCount和EventTakeSuccessCount这样的字段。前者是 Source 放进 Channel 的累计条数后者是 Sink 从 Channel 取走并成功发到 Kafka 的累计条数。两个数字都持续增长就说明整条链路健康。确认一切正常后再 CtrlC 停掉前台进程改用 nohup 后台运行nohup /opt/module/flume/bin/flume-ng agent \ --name a1 \ --conf /opt/module/flume/conf \ --conf-file /opt/module/flume/conf/flume_kafka.conf \ -Dflume.root.loggerINFO,LOGFILE \ 21 注意后台启动日志就不再打到控制台了默认写入 Flume 安装目录下的 logs/flume.log后续排障要看这个文件。3.4 多文件分组和简单过滤的顺手优化这一步不是必须的但真实项目里很有用。Taildir Source 支持同时监控多个文件组如果你不仅有应用日志还有 Nginx 访问日志可以把配置扩展成a1.sources.r1.filegroups f1 f2 a1.sources.r1.filegroups.f1 /opt/module/applog/log/app.*.log a1.sources.r1.filegroups.f2 /opt/module/applog/log/nginx.*.log a1.sources.r1.basenameHeader truebasenameHeader true会在 Event 的 header 里记录文件的基本名称这样下游 Flink 拿到数据时可以方便地区分这条日志来自应用服务还是 Nginx。还有一个很常见的需求是过滤掉一些没用的日志行。比如探活接口每秒钟都会刷一条 health_check 日志你不想让这种东西进 Kafka就可以加一个 interceptora1.sources.r1.interceptors i1 a1.sources.r1.interceptors.i1.type regex_filter a1.sources.r1.interceptors.i1.regex health_check a1.sources.r1.interceptors.i1.excludeEvents true这段配置的意思是凡是不包含health_check这个关键字的行都保留包含的直接丢弃。这类小技巧在课程笔记里不一定细讲但实际工作中几乎天天用提前掌握能省很多麻烦。4. 排障实录这节笔记遇到的坑和排查套路4.1 最经典的报错找不到 KafkaProducer 类这一节的第一个大坑也是我记忆最深的。启动 Flume 时日志里出现了ClassNotFoundException: org.apache.kafka.clients.producer.KafkaProducer直接导致 Kafka Sink 初始化失败。原因很明确Flume 自身带的 Kafka 客户端库版本太旧而 Kafka 集群版本是 2.x旧客户端根本演不了这个协议。解决办法是把 Kafka 安装目录下的 kafka-clients jar 包拷贝到 Flume 的 lib 目录下cp /opt/module/kafka/libs/kafka-clients-2.4.1.jar /opt/module/flume/lib/然后重启 Flume。文件名后面的版本号要根据你 Kafka 的实际版本写不过它们的命名规律是一致的。这一点我要特别强调不要为了“保险”把 Kafka 的 libs 目录整个拷到 Flume 的 lib 里没必要还会引发 jar 包冲突。kafka-clients 一个包就够用了它负责客户端和 broker 之间的协议通信、序列化、网络传输Kafka Sink 依赖的核心就是它。如果重启之后 Flume 日志里还报类似 Kafka producer 配置不识别之类的错先检查你拷贝的 jar 包版本是不是跟 Kafka broker 差太多。版本差距过大时客户端虽然能连接但可能出现请求超时或协议协商失败这类问题最隐蔽往往让你在 Kafka 集群上排查半天实际上就是 jar 包版本不匹配。4.2 消费不到数据从哪儿查起另一类常见问题是Flume 启动没报错日志文件也在增长但 Kafka 消费者那边就是看不到数据。这种问题我建议按数据流向顺序排查别一上来就怀疑 Kafka。第一步确认日志文件真的有新内容。用tail -f直接看文件如果文件都没增长那问题根本不存在。第二步确认 Flume 进程还活着ps -ef | grep flume第三步看 Flume 日志里的计数器重点是 EventPutSuccessCount。如果这个数字一直不变说明 Source 没读到文件。这时候多半是 filegroups 的路径正则没匹配上。我当时犯过一个低级错误日志实际文件名是app.2024-01-01.log我配置写成了app.log没有带通配符Flume 监控不到自然一条数据都读不到。去taildir_position.json里看有没有出现你目标文件的绝对路径一目了然。第四步如果 put 在涨、take 不涨说明 Channel 里积压了数据问题是 Kafka Sink 发不出去。这时候用 Atlassian 的 telnet 或者 nc 测一下 broker 端口通不通确认 bootstrap.servers 里写的地址能从 Flume 所在机器访问。集群环境下最容易出现的问题就是只写了hadoop102:9092但该节点的 broker 恰好没起来连不上。最后再确认 Kafka topic 名称写没写对。Kafka Sink 配置里kafka.topic拼错了消息会被发到不存在的主题上消费者自然收不到。但注意Kafka 默认允许自动创建主题你拼错也会被自动建出来所以只看 topic 存不存在还不够要看它里面到底有没有数据。用 console-consumer 从 beginning 消费时如果没人消费过的 topic 也显示空极有可能是消息发到了别的主题里。4.3 数据重复和丢失的经典局面再说两个跟可靠性相关的坑。重复消费最常见的触发方式是手贱删了 positionFile。positionFile 一旦不存在Flume 重启后会扫描目标目录里的所有文件从文件头开始重新读一遍。数据量小的时候看不出问题几 GB 的日志文件重新灌进 Kafka下游实时计算马上就会出现重复统计。我当时为了测试还真的删过 positionFile删完就后悔了因为 Kafka 里瞬间多出一大堆重复数据。所以除非你真的想让 Flume 从头读取历史文件否则永远不要动 positionFile。丢失数据Memory Channel 的可靠性就是这样Flume 进程突然崩溃Channel 里还没发往 Kafka 的数据就没了。我建议做两件事防丢。第一给 Kafka Sink 配acks 1或者-1不要配成 0。acks0意味着 producer 不管 broker 收到没有直接认为发送成功这种配置下 Kafka 侧一个闪断就会造成大量消息丢失。第二如果业务要求高可靠换用 File Channel。具体配置就是把a1.channels.c1.type改成file再加上 dataDirs 指向一个持久化目录。File Channel 的吞吐确实低一些但数据不会因为进程崩溃就蒸发了这属于可靠性和性能的取舍学习阶段了解即可。还有一个容易被忽略的重复场景日志文件被切割。日志系统按天滚动文件时app.2024-01-01.log被 rename 成归档文件同时新建一个同名文件。TaildirSource 用 inode 识别文件rename 后老文件还在被 Flume 记录新文件又产生新的 inode如果两边内容有重叠就可能读到重复数据。对这种场景我的经验是让日志框架在写入时保持文件基本名稳定不要在同一天内反复创建同名文件。4.4 高频问题速查表整理一个速查表方便你以后直接对照现象可能原因排查方向 / 对策启动报 ClassNotFoundExceptionFlume 自带 kafka-clients 版本过旧拷贝 Kafka libs 下的 kafka-clients jar 到 Flume lib启动报 Agent configuration 语法错误组件名或属性名拼写不对核对 sources/channels/sinks 的三段绑定关系日志文件一直在涨但 Source 不读filegroups 正则没匹配到文件名检查 taildir_position.json 里有没有对应文件路径Source 在涨但 Sink 不涨Kafka Sink 连不上 brokertelnet 测端口、确认 bootstrap.servers、查看 Flume 日志有数据进入 Kafka 但消费端看不到topic 名称不一致或 group offset 问题确认 kafka.topic 拼写、用 --from-beginning 消费删了 positionFile 后重复消费Flume 从头读取所有文件不要随意删除 positionFileChannel full 报错capacity 太小导致反压调大 capacity 和 transactionCapacity日志权限不足Flume 进程用户读不了日志文件chown 或 chmod 调整文件权限这些坑我后来在好几个项目里都碰到过尤其是 kafka-clients 版本那个报错几乎每次从零搭环境都要来一遍。配好一次之后把常用的 jar 包和配置模板存下来能省掉后面不少无意义的排障时间。我当时学到这里最大的体会是Flume 到 Kafka 的配置本身不难难在理解每个参数为什么这样设以及遇到问题后知道往哪查。你要是能把flume_kafka.conf这份配置自己手写一遍再把日志文件 echo 进目录最后在 Kafka console consumer 里看到消息那这一节就算真正拿下了。接下来无论是去连 Flink 做实时计算还是把数据落 HDFS 进离线数仓你都会感谢这节内容打下的底子。
返回列表