Flume-ng 拖尾文件


我试图了解如何使用 Flume-ng 尾部文件,以便可以将数据推送到 HDFS 中。在第一个实例中,我设置了一个简单的conf文件:

tail1.sources = source1
tail1.sinks = sink1
tail1.channels = channel1

tail1.sources.source1.type = exec
tail1.sources.source1.command = tail -F /var/log/apache2/access.log
tail1.sources.source1.channels = channel1

tail1.sinks.sink1.type = logger

tail1.channels.channel1.type = memory
tail1.channels.channel1.capacity = 1000
tail1.channels.channel1.transactionCapacity = 100

tail1.sources.source1.channels = channel1
tail1.sinks.sink1.channel = channel1


flume-ng agent --conf-file tail1.conf -n tail1 -Dflume.root.logger=DEBUG,INFO,console


12/12/05 11:01:07 信息生命周期.LifecycleSupervisor: 开始 生命周期主管 1 12/12/05 11:01:07 INFO 节点.FlumeNode: Flume 节点开始 - tail1 12/12/05 11:01:07 INFO nodemanager.DefaultLogicalNodeManager:节点管理器从 2005 年 12 月 12 日开始 11:01:07 信息生命周期.LifecycleSupervisor:开始生命周期 主管 8 12/12/05 11:01:07 信息 properties.PropertiesFileConfigurationProvider:配置提供程序 开始于 2005 年 12 月 12 日 11:01:07 信息 properties.PropertiesFileConfigurationProvider:正在重新加载 配置文件:tail1.conf 12/12/05 11:01:07 INFO conf.FlumeConfiguration:处理:sink1 12/12/05 11:01:07 信息 conf.FlumeConfiguration:处理:sink1 12/12/05 11:01:07 信息 conf.FlumeConfiguration:添加接收器:sink1 代理:tail1 12/12/05 11:01:07 INFO conf.FlumeConfiguration:验证后水槽 配置包含代理配置:[tail1] 12/12/05 11:01:07 信息属性.PropertiesFileConfigurationProvider:创建 通道 12/12/05 11:01:08 INFO 仪器。MonitoredCounterGroup: 监控的计数器组类型:CHANNEL,名称:channel1,已注册 成功地。 2005 年 12 月 12 日 11:01:08 信息 properties.PropertiesFileConfigurationProvider:创建的通道 通道1 12/12/05 11:01:08 信息接收器.DefaultSinkFactory:正在创建 接收器实例:sink1,类型:记录器 12/12/05 11:01:08 INFO nodemanager.DefaultLogicalNodeManager:开始新配置:{ sourceRunners:{source1=EventDrivenSourceRunner: { 来源:org.apache.flume.source.ExecSource@1839aa9 }} inkRunners:{sink1=SinkRunner:{ 策略:org.apache.flume.sink.DefaultSinkProcessor@11f0c98 counterGroup:{ name:null 计数器:{} } }} 频道:{channel1=org.apache.flume.channel.MemoryChannel@17​​40f55} } 12/12/05 11:01:08 信息 nodemanager.DefaultLogicalNodeManager:正在启动 频道频道1 12/12/05 11:01:08 信息 Instrumentation.MonitoredCounterGroup:组件类型:CHANNEL,名称: 频道 1 于 2005 年 12 月 12 日 11:01:08 开始 nodemanager.DefaultLogicalNodeManager:启动接收器 sink1 12/12/05 11:01:08 INFO nodemanager.DefaultLogicalNodeManager:启动源 source1 12/12/05 11:01:08 信息 source.ExecSource:Exec 源启动 使用命令:tail -F /var/log/apache2/access.log



tail -F /var/log/apache2/access.log

我可以在哪里看到正在写入的文件: - - [05/Dec/2012:10:58:07 +0000] "GET / HTTP/1.1" 200 483 "-" "Mozilla/5.0 (Windows NT 6.1; WOW64) AppleWebKit/537.11 (KHTML, like Gecko) Chrome/23.0.1271.95 Safari/537.11" - - [05/Dec/2012:10:58:07 +0000] "GET /favicon.ico HTTP/1.1" 404 502 "-" "Mozilla/5.0 (Windows NT 6.1; WOW64) AppleWebKit/537.11 (KHTML, like Gecko) Chrome/23.0.1271.95 Safari/537.11" - - [05/Dec/2012:10:58:21 +0000] "GET / HTTP/1.1" 304 209 "-" "Mozilla/5.0 (Windows NT 6.1; WOW64) AppleWebKit/537.11 (KHTML, like Gecko) Chrome/23.0.1271.95 Safari/537.11" - - [05/Dec/2012:10:58:22 +0000] "GET /favicon.ico HTTP/1.1" 404 502 "-" "Mozilla/5.0 (Windows NT 6.1; WOW64) AppleWebKit/537.11 (KHTML, like Gecko) Chrome/23.0.1271.95 Safari/537.11"


你的配置文件看起来不错。我在 CDH4 中使用它并按照您的预期工作,我所做的只是更改尾部的日志文件位置。我看到控制台上的输出。就我而言,新的日志数据被连续写入我跟踪的文件中。数据中的时间戳使得您的示例中的情况看起来可能并非如此。

这是一个更完整的conf示例,更符合我认为您想要做的事情。它将跟踪文件并每 10 分钟或 10K 记录写入一个新的 HDFS 文件。改变agent1.sources.source1.命令到你的 tail 命令并改变agent1.sinks.sink1.hdfs.path and agent1.sinks.sink1.hdfs.filePrefix基于您的 HDFS 配置。

# A single-node Flume configuration
# uses exec and tail and will write a file every 10K records or every 10 min
# Name the components on this agent
agent1.sources = source1
agent1.sinks = sink1
agent1.channels = channel1

# Describe/configure source1
agent1.sources.source1.type = exec
agent1.sources.source1.command = tail -f /home/cloudera/LogCreator/fortune_log.log

# Describe sink1
agent1.sinks.sink1.type = hdfs
agent1.sinks.sink1.hdfs.path = hdfs://localhost/flume/logtest/
agent1.sinks.sink1.hdfs.filePrefix = LogCreateTest
# Number of seconds to wait before rolling current file (0 = never roll based on time interval)
agent1.sinks.sink1.hdfs.rollInterval = 600
# File size to trigger roll, in bytes (0: never roll based on file size) 
agent1.sinks.sink1.hdfs.rollSize = 0
#Number of events written to file before it rolled (0 = never roll based on number of events) 
agent1.sinks.sink1.hdfs.rollCount = 10000
# number of events written to file before it flushed to HDFS 
agent1.sinks.sink1.hdfs.batchSize = 10000 
agent1.sinks.sink1.hdfs.txnEventMax = 40000
# -- Compression codec. one of following : gzip, bzip2, lzo, snappy
# hdfs.codeC = gzip
#format: currently SequenceFile, DataStream or CompressedStream
#(1)DataStream will not compress output file and please don't set codeC
#(2)CompressedStream requires set hdfs.codeC with an available codeC
agent1.sinks.sink1.hdfs.fileType = DataStream 
# -- "Text" or "Writable"
agent1.sinks.sink1.hdfs.appendTimeout = 10000
agent1.sinks.sink1.hdfs.callTimeout = 10000
# Number of threads per HDFS sink for HDFS IO ops (open, write, etc.)
# Number of threads per HDFS sink for scheduling timed file rolling
agent1.sinks.sink1.hdfs.rollTimerPoolSize = 1 
# hdfs.kerberosPrin--cipal Kerberos user principal for accessing secure HDFS
# hdfs.kerberosKey--tab Kerberos keytab for accessing secure HDFS
# hdfs.round false Should the timestamp be rounded down (if true, affects all time based escape sequences except %t)
# hdfs.roundValue1 Rounded down to the highest multiple of this (in the unit configured using
# hdfs.roundUnit), less than current time.
# hdfs.roundUnit second The unit of the round down value - second, minute or hour.
# serializer TEXT Other possible options include AVRO_EVENT or the fully-qualified class name of an implementation of the EventSerializer.Builder interface.
# serializer.*

# Use a channel which buffers events to a file
# -- The component type name, needs to be FILE.
agent1.channels.channel1.type = FILE 
# checkpointDir ~/.flume/file-channel/checkpoint The directory where checkpoint file will be stored
# dataDirs ~/.flume/file-channel/data The directory where log files will be stored
# The maximum size of transaction supported by the channel
agent1.channels.channel1.transactionCapacity = 1000000 
# Amount of time (in millis) between checkpoints
agent1.channels.channel1.checkpointInterval 30000
# Max size (in bytes) of a single log file 
agent1.channels.channel1.maxFileSize = 2146435071
# Maximum capacity of the channel 
agent1.channels.channel1.capacity 10000000 
#keep-alive 3 Amount of time (in sec) to wait for a put operation
#write-timeout 3 Amount of time (in sec) to wait for a write operation

# Bind the source and sink to the channel
agent1.sources.source1.channels = channel1
agent1.sinks.sink1.channel = channel1

