2.2.5 实时监控端口数据,输出至HDFS 1)案例需求:使用Flume监听整个目录的文件,并上传至HDFS和本地logger 2)实现步骤: 使用Flume监听一个端口,收集该端口数据,并输出到hdfs [root@kb129 myconf2]# vim ./demo.conf # Describe/configure the source a5.sources.r1.type=netcat a5.sources.r1.bind=localhost # Use a channel which buffers events in memory a5.channels.c1.type=memory a5.channels.c1.capacity=1000 a5.channels.c1.transactionCapacity=100 a5.sinks.k1.hdfs.fileType=DataStream a5.sinks.k1.hdfs.filePrefix=flumetohdfs a5.sinks.k1.hdfs.fileSuffix=.txt a5.sinks.k1.hdfs.path=hdfs://kb129:9000/kb23flume2/ # Bind the source and sink to the channel a5.sources.r1.channels=c1
启动监听 [root@kb129 flume190]# ./bin/flume-ng agent -n a5 -c ./conf/ -f ./conf/myconf2/demo.conf -Dflume.root.logger=INFO,console 开启通信端口,发送数据,在hdfs查看 2.2.6 实时监控目录下指定文件(spooldir、正则、断点续传、sink至kafka) Exec source适用于监控一个实时追加的文件,不能实现断点续传;Spooldir Source适合用于同步新文件,但不适合对实时追加日志的文件进行监听并同步;而Taildir Source适合用于监听多个实时追加的文件,并且能够实现断点续传。总结:exec source从外部命令的输出中读取数据,spooldir source适用于处理已存在的文件,而taildir source适用于实时产生的日志文件。 (1)读取文件,管道类型memory [root@kb129 myconf2]# vim events-flume-logger.conf events.sources=eventsSource events.channels=eventsChannel # Describe/configure the source events.sources.eventsSource.type=spooldir events.sources.eventsSource.spoolDir=/opt/kb23/flumelogfile/events events.sources.eventsSource.deserializer=LINE events.sources.eventsSource.deserializer.maxLineLength=32000 events.sources.eventsSource.includePattern=events.csv # Use a channel which buffers events in memory events.channels.eventsChannel.type=memory events.sinks.eventsSink.type=logger # Bind the source and sink to the channel events.sources.eventsSource.channels= eventsChannel events.sinks.eventsSink.channel= eventsChannel
(2)读取文件,加正则过滤,管道类型file 1. 文件管道(file):文件管道将事件数据写入磁盘上的文件中。这种管道适用于需要持久化存储数据的场景。文件管道可以在系统重启后继续读取未处理的数据,因为数据被写入了磁盘。但是,由于写入和读取都需要磁盘操作,所以文件管道的性能可能相对较低。 2. 内存管道(memory):内存管道将事件数据保存在内存中。这种管道适用于对性能要求较高的场景,因为内存操作比磁盘操作更快。内存管道不会将数据持久化到磁盘,所以在系统重启后,未处理的数据将会丢失。因此,内存管道适用于对数据可靠性要求不高的场景。 [root@kb129 myconf2]# vim events-flume-logger.conf events.sources=eventsSource events.channels=eventsChannel # Describe/configure the source events.sources.eventsSource.type=spooldir events.sources.eventsSource.spoolDir=/opt/kb23/flumelogfile/events events.sources.eventsSource.deserializer=LINE events.sources.eventsSource.deserializer.maxLineLength=32000 events.sources.eventsSource.includePattern=events_[0-9]{4}-[0-9]{2}-[0-9]{2}.csv events.sources.eventsSource.interceptors=head_filter events.sources.eventsSource.interceptors.head_filter.type=regex_filter events.sources.eventsSource.interceptors.head_filter.regex=^event_id* events.sources.eventsSource.interceptors.head_filter.excludeEvents=true # Use a channel which buffers events in memory events.channels.eventsChannel.type=file events.channels.eventsChannel.checkpointDir=/opt/kb23/checkpoint/events events.channels.eventsChannel.dataDirs=/opt/kb23/data/events events.sinks.eventsSink.type=logger # Bind the source and sink to the channel events.sources.eventsSource.channels=eventsChannel events.sinks.eventsSink.channel=eventsChannel
 (3)sink输出指向kafka,topic [root@kb129 myconf2]# vim ./events-flume-kafka.conf events.sources=eventsSource events.channels=eventsChannel # Describe/configure the source events.sources.eventsSource.type=spooldir events.sources.eventsSource.spoolDir=/opt/kb23/flumelogfile/events events.sources.eventsSource.deserializer=LINE events.sources.eventsSource.deserializer.maxLineLength=32000 events.sources.eventsSource.includePattern=events_[0-9]{4}-[0-9]{2}-[0-9]{2}.csv events.sources.eventsSource.interceptors=head_filter events.sources.eventsSource.interceptors.head_filter.type=regex_filter events.sources.eventsSource.interceptors.head_filter.regex=^event_id* events.sources.eventsSource.interceptors.head_filter.excludeEvents=true # Use a channel which buffers events in memory events.channels.eventsChannel.type=file events.channels.eventsChannel.checkpointDir=/opt/kb23/checkpoint/events events.channels.eventsChannel.dataDirs=/opt/kb23/data/events events.sinks.eventsSink.type=org.apache.flume.sink.kafka.KafkaSink events.sinks.eventsSink.topic=events events.sinks.eventsSink.brokerList=192.168.142.129:9092 events.sinks.eventsSink.batchSize=640 # Bind the source and sink to the channel events.sources.eventsSource.channels=eventsChannel events.sinks.eventsSink.channel=eventsChannel
 |