请描述在 Apache Flume 中实现自定义 Source 以采集特定格式数据的具体方法,包括需要继承的类、实现的核心接口以及配置方式。
考察说明
考察对 Flume 自定义 Source 机制的理解,包括继承关系、生命周期方法和配置加载。
回答思路
- 【回答框架 1】Flume 中的 Source 负责从外部系统接收数据并写入 Channel。自定义 Source 通常继承 AbstractSource 类,该类提供了 ChannelProcessor 的引用,并实现 Configurable 和 PollableSource 或 EventDrivenSource 接口。
- 【回答框架 2】对于 EventDrivenSource,需要重写 start() 和 stop() 方法,在 start() 中启动事件生成线程,持续调用 getChannelProcessor().processEvent(event) 或 processEventBatch(batch) 将事件写入 Channel。
- 【回答框架 3】对于 PollableSource,需要实现 process() 方法,该方法返回 Status.READY 或 Status.BACKOFF,并负责获取数据并写入 Channel。
- 【回答框架 4】实现 Configurable 接口的 configure(Context context) 方法,用于读取配置文件中的自定义参数,例如通过 context.getString("param", "default") 获取配置值。
- 【回答框架 5】在构建 Flume 配置时,通过 type 指定自定义 Source 类的全限定名,并通过 prefix 传递自定义参数,例如 agent.sources.r1.type = com.example.MySource。
- 【关键点 1】自定义 Source 必须继承 AbstractSource 并实现 Configurable 接口。
- 【关键点 2】根据数据采集方式选择实现 EventDrivenSource 或 PollableSource 接口。
- 【关键点 3】重写 configure() 方法以读取外部配置。
- 【关键点 4】最终通过 getChannelProcessor().processEvent() 或 processEventBatch() 将数据写入 Channel。
- 【易错点 1】不要在 start() 中阻塞主线程,否则会导致 Flume 无法正常启动。
- 【易错点 2】确保在 stop() 中释放所有资源,防止内存泄漏。
- 【易错点 3】配置参数名称需与 configure() 中读取的键一致,否则会使用默认值。