flume扩展
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。