小菜鸟

java菜鸟号正在起航

Flume 内置组件不够用?自己写一个,没你想的那么难

Flume 内置的 Source、Sink、拦截器覆盖了大部分场景,但总有例外。

我遇到过的情况:业务方发的日志格式是自定义的二进制协议,内置的 Source 根本不认识;下游存储是个自研的时序数据库,没有现成的 Sink。这时候只能自己写。

说实话,Flume 的扩展接口设计得还算清晰,继承几个类、实现几个方法就能跑起来。这篇文章把三种扩展——拦截器、Source、Sink——的写法都过一遍,代码贴全,配置也附上,你需要的时候直接抄。

开发前准备:开发环境与依赖

Flume 的自定义组件本质上就是写 Java 类,打包成 JAR 扔到 Flume 的 lib 目录下。

依赖配置

pom.xml 中添加 Flume 核心依赖(以 1.9.0 为例):

<dependency>  
    <groupId>org.apache.flume</groupId>  
    <artifactId>flume-ng-core</artifactId>  
    <version>1.9.0</version>  
    <scope>provided</scope>  <!-- 运行时由 Flume 环境提供 -->  
</dependency>

scope 设成 provided,因为 Flume 运行时已经带了这些类,打包的时候别打进去,避免版本冲突。

核心接口

Flume 扩展的核心是实现官方定义的接口,各组件对应的接口如下:

组件类型 需实现的接口 / 继承的类 核心方法
拦截器 org.apache.flume.interceptor.Interceptor intercept(Event) 处理单个事件
Source 继承 AbstractSource,实现 PollableSource process() 产生并发送事件
Sink 继承 AbstractSink,实现 Configurable process() 从 Channel 消费事件

每个组件都要配一个 Builder 类,Flume 通过 Builder 来实例化你的组件。这块容易忘,后面代码里会看到。

实战一:自定义拦截器(Interceptor):给 Event 打个标签

拦截器的用处是:数据从 Source 出来、进 Channel 之前,你插一手,改改 Header、改改 Body,或者直接丢掉某些数据。

下面这个例子实现了一个简单的分类器:根据 Body 内容给 Event 的 Header 里加一个 type 字段,后续可以用 Multiplexing Channel Selector 按类型路由到不同的 Channel。

阅读全文 »

Flume Sink 挂了怎么办?接收处理器帮你兜底

Flume 跑在生产上,最怕两件事:Sink 挂了数据发不出去,和单个 Sink 扛不住流量积压

前几篇文章里,每个 Agent 都只配了一个 Sink。单点跑的没问题,但线上环境谁敢用单点?Sink 指向的 Kafka 集群抖动、HDFS 做安全模式、网络闪断,随便一个就能让采集链路瘫痪。

Flume 的 Sink Group + 接收处理器(Processor) 就是干这个的。这篇文章我讲清楚两种处理器怎么配、什么场景用哪个,以及我实际用下来的一些判断。

先搞清楚处理器是干什么的

一个 Agent 可以有多个 Sink,多个 Sink 组成一个 Sink Group。处理器(Processor)管的是:数据来了,往哪个 Sink 发?

Flume 官方提供三种接收处理器:

处理器类型 核心功能 适用场景
DefaultSinkProcessor 单 Sink 处理(不支持组) 简单场景,无需冗余或负载均衡
FailoverSinkProcessor 故障转移(按优先级切换)主 Sink 挂了切到备 Sink 需要高可用性的关键链路
LoadBalancingSinkProcessor 负载均衡(轮询或随机)多个 Sink 分摊流量 需要提升吞吐量的高并发场景
阅读全文 »

flume事务机制:不丢数据不重数据,靠的就是这两件事

做数据采集,最怕两件事:数据丢了数据重了

丢了,下游报表对不上;重了,统计结果全偏。Flume 能在大数据场景下保证数据的准确性,靠的就是它的事务机制

说白了,Flume 在数据流转的两个关键节点——Source 往 Channel 写、Sink 从 Channel 读——都加了事务控制。写成功了才算数,写失败了就回滚重来。这篇文章我把这个机制拆开讲清楚,顺便说几个实际场景里的事务问题。

为什么需要事务?

Flume 在线上跑着,会遇到各种意外:

  • Kafka 集群抖动,Source 写不进去
  • HDFS 满了,Sink 写不动
  • Flume 进程本身挂了,重启

如果没有事务,数据写到一半出问题,就可能出现”Channel 里记了一笔但实际没发出去”或者”发出去了但 Channel 没删掉”这种状态不一致的情况。

事务要解决的就是这个问题:要么全成,要么全不成,没有中间状态

Flume 事务的两大阶段

Flume 的事务机制贯穿数据流转的全流程,分为Put 事务(Source → Channel)和Take 事务(Channel → Sink),两个阶段独立保障数据可靠性。

阅读全文 »

Kafka 里的数据怎么落到 HDFS?Flume 这套配置我跑通了

Kafka 和 HDFS 是大数据里最常见的两个存储,一个管实时、一个管离线。中间怎么打通?Flume 就是干这个的。

方案架构:Kafka → Flume → HDFS 的数据流转

本方案的核心是通过 Flume 构建 “Kafka 消费 → 通道缓冲 → HDFS 写入” 的流水线,具体架构如下:

注意这里有个特殊的地方:Channel 用的是 KafkaChannel,不是 Memory 或 File。也就是说,Flume 中间缓冲的数据也放在 Kafka 里,利用 Kafka 自身的持久化保证可靠性。这个方案的优点是可靠,缺点是得多维护一个 Kafka topic。

graph LR
    A[Kafka 主题 kafka-flume] -->|Kafka Source 消费| B[Flume Agent]
    B -->|Kafka Channel 缓冲| C[Kafka 通道主题 kafka-channel]
    B -->|HDFS Sink 写入| D[HDFS 存储 按时间分区]
  • Kafka Source:从 Kafka 主题(如 kafka-flume)消费消息;
  • Kafka Channel:使用 Kafka 作为 Flume 的通道,利用 Kafka 的可靠性暂存数据;
  • HDFS Sink:将数据按时间分区(如 %Y-%m-%d/%H)写入 HDFS,便于后续离线分析。

详细配置步骤

阅读全文 »

flume监控日志写到 Kafka,我为什么不用 log4j2 直连

日志采集有一个常见的争议:应用到底该直接把日志打到 Kafka,还是先写本地文件再让 Flume 采?

我之前也纠结过这个问题。后来线上出过一次事故——Kafka 集群抖动,log4j2 的 Kafka appender 直接把应用线程堵死了,服务全挂。从那以后我就认了:应用只管写文件,采集的事交给 Flume

这篇我就讲一下 Flume 监控文件写到 Kafka 的配置方法,以及几个我用下来觉得关键的参数。

为什么中间要加一层 Flume?

直接让应用通过 log4j2 写 Kafka,确实省事,配置几行就搞定。但带来的问题是:

  • Kafka 出问题会影响应用。Broker 不可用、网络抖动、topic 被删,应用的生产者线程会阻塞,业务接口跟着超时。
  • 应用要关心 Kafka 的配置。bootstrap.servers、序列化方式、acks 这些参数都塞到应用的配置文件里,跟业务配置混在一起。
  • 没法做数据清洗。log4j2 的 appender 只负责发,你想在写入前做格式转换、过滤敏感信息,基本没法搞。

加了 Flume 这一层之后:

  • 应用只写本地文件,write 完就结束,Kafka 死活跟它没关系
  • Flume 负责采、负责发、负责重试,采集挂了也不影响业务
  • 想清洗数据?加个拦截器就行,随时改随时生效
阅读全文 »
0%