Flume:Flume在电商数仓中的应用

Flume用于采集日志数据。
采集Flume

1、Source
选择Taildir Source
Taildir Source 是 Apache Flume 中的一种数据源(Source)类型,用于监控一个或多个日志文件的新增内容,并将这些新增内容作为事件(Event)发送到 Flume 的 Channel 中。它是 Exec Source 和 Spooling Directory Source 的更高效、更可靠的替代方案。
(1)工作原理
Taildir Source 会持续跟踪指定日志文件的变化,通过记录文件的偏移量(offset)信息,只读取文件中新增的内容。它使用一个位置文件(通常是 JSON 格式)来存储每个被监控文件的当前读取位置,这样即使 Flume 进程重启,也能从上次中断的位置继续读取文件内容,实现断点续传(Apache 1.7)。
(2)主要特点
- 断点续传:如上述所说,Taildir Source 可以记录每个文件的读取偏移量,在 Flume 重启后,能够自动从上次中断的位置继续读取文件,避免数据丢失和重复读取。
- 多文件监控:可以同时监控多个日志文件,并且能够高效地处理大量文件的变化。例如,在一个大型服务器集群中,每个服务器可能会产生多个日志文件,Taildir Source 可以同时监控这些文件,将日志数据收集到 Flume 中。
- 性能高效:相比传统的
Exec Source(通过执行外部命令来读取文件),Taildir Source 不需要频繁地启动新的进程,减少了系统开销,提高了性能。同时,它采用了异步 I/O 操作,能够更高效地处理文件的读取。 - 动态文件发现:支持动态发现新的日志文件。当有新的日志文件被创建时,Taildir Source 可以自动检测到并开始监控该文件,无需手动重启 Flume 进程。
(3)应用场景
Taildir Source 主要用于日志收集场景,例如:
- 服务器日志收集:收集 Web 服务器(如 Apache、Nginx)、应用服务器(如 Tomcat、Jetty)等产生的日志文件,将这些日志数据传输到大数据平台(如 Hadoop、Elasticsearch)进行存储和分析。
- 系统日志监控:监控操作系统产生的各种日志文件(如
/var/log/syslog、/var/log/messages等),及时发现系统中的异常事件和错误信息。
(4)如何处理重复数据
客户端、埋点等可能多次上报,造成数据重复。通常不处理,下游dwd层group by、窗口函数排序取最新一条
2、Channel
选择的是kafka channel。数据存储在kafka,基于磁盘,可靠性高,传输速度也快,省去了sink阶段。
3、Sink
kafka channel省去了sink阶段
kafka topic设置3个,存放错误日志主题、启动日志主题、时间日志主题。
4、自定义Flume拦截器
(1)ETL拦截器:过滤时间戳不合法和Json数据不完整的日志
(2)日志类型区分拦截器:启动日志和事件日志区分开来,方便发往Kafka的不同Topic
取FLume接收消息头header,从json数据中的eventType取出响应的日志类型,将日志类型存储到header中实现不同事件发送到不同topic
(3)自定义拦截器步骤

(4)优点、缺点
优点:数据只处理1次,轻度处理。
缺点:影响传输效率,不适合对实时性要求高的数据处理。
没有拦截器也可以,对于重复数据,下游dwd层group by、窗口函数排序取最新一条
消费Flume

1、Source
选择kafka source,因为上游是kafka topic。
2、Channel
选择的是file channel.
3、Sink
选择的是hdfs sink。
4、自定义Flume时间戳拦截器
由于Flume默认会用Linux系统时间,作为输出到HDFS路径的时间。如果数据是23:59分产生的。Flume消费Kafka里面的数据时,有可能已经是第二天了,那么这部门数据会被发往第二天的HDFS路径。我们希望的是根据日志里面的实际时间,发往HDFS的路径,所以下面拦截器作用是获取日志中的实际时间。
解决的思路:
拦截json日志,通过fastjson框架解析json,获取实际时间ts。将获取的ts时间写入拦截器header头,header的key必须是timestamp,因为Flume框架会根据这个key的值识别为时间,写入到HDFS。
更多推荐

所有评论(0)