Flume 1.6.0 Spooling Directory Source添加事件头Timestamp问题咨询
解决Flume事件头缺失timestamp的问题
你遇到的Expected timestamp in the Flume event headers, but it was null错误,是因为你的HDFS Sink依赖事件Header中的timestamp字段来生成配置里的时间分区路径(/HADOOP/PATH/%Y/%m/%d/%H/),但当前SpoolDir Source并没有把JSON文件里的timestamp字段同步到事件Header中。下面是具体的解决方案:
核心思路:用拦截器提取JSON中的timestamp到事件Header
Flume提供了RegexExtractorInterceptor可以从事件内容中提取指定字段,再结合TimestampInterceptor(如果需要格式转换),就能把JSON里的timestamp放到事件Header里。
情况1:JSON中的timestamp是毫秒级数字(如"timestamp": 1716191400000)
直接添加正则提取拦截器即可,修改你的Source配置部分:
agent.sources = file agent.channels = channel agent.sinks = hdfsSink # SOURCES CONFIGURATION - 新增拦截器配置 agent.sources.file.type = spooldir agent.sources.file.channels = channel agent.sources.file.spoolDir = /path/to/json_files # 添加拦截器 agent.sources.file.interceptors = tsInterceptor agent.sources.file.interceptors.tsInterceptor.type = regex_extractor # 正则匹配JSON中的数字型timestamp,注意转义特殊字符 agent.sources.file.interceptors.tsInterceptor.regex = "timestamp":\\s*(\\d+) # 将提取到的值存入名为timestamp的Header agent.sources.file.interceptors.tsInterceptor.serializers = ts agent.sources.file.interceptors.tsInterceptor.serializers.ts.name = timestamp # SINKS CONFIGURATION(保持原配置不变) agent.sinks.hdfsSink.type = hdfs agent.sinks.hdfsSink.hdfs.path = /HADOOP/PATH/%Y/%m/%d/%H/ agent.sinks.hdfsSink.hdfs.filePrefix = common agent.sinks.hdfsSink.hdfs.fileSuffix = .json agent.sinks.hdfsSink.hdfs.rollInterval = 300 agent.sinks.hdfsSink.hdfs.rollSize = 5242880 agent.sinks.hdfsSink.hdfs.rollCount = 0 agent.sinks.hdfsSink.hdfs.maxOpenFiles = 2 agent.sinks.hdfsSink.hdfs.fileType = DataStream agent.sinks.hdfsSink.hdfs.callTimeout = 100000 agent.sinks.hdfsSink.hdfs.batchSize = 1000 agent.sinks.hdfsSink.channel = channel # CHANNELS CONFIGURATION(保持原配置不变) agent.channels.channel.type = memory agent.channels.channel.capacity = 10000 agent.channels.channel.transactionCapacity = 1000
情况2:JSON中的timestamp是字符串格式(如"timestamp": "2024-05-20 14:30:00")
这种情况需要先提取字符串,再转换成Flume需要的毫秒级时间戳,需要两个拦截器配合:
agent.sources = file agent.channels = channel agent.sinks = hdfsSink # SOURCES CONFIGURATION - 新增两个拦截器 agent.sources.file.type = spooldir agent.sources.file.channels = channel agent.sources.file.spoolDir = /path/to/json_files # 拦截器顺序:先提取字符串,再转换为时间戳 agent.sources.file.interceptors = tsExtractor tsConverter # 第一个拦截器:从JSON中提取字符串型timestamp agent.sources.file.interceptors.tsExtractor.type = regex_extractor agent.sources.file.interceptors.tsExtractor.regex = "timestamp":\\s*"([^"]+)" agent.sources.file.interceptors.tsExtractor.serializers = tsStr agent.sources.file.interceptors.tsExtractor.serializers.tsStr.name = timestamp_str # 第二个拦截器:将字符串转换为毫秒时间戳,注意匹配你的时间格式 agent.sources.file.interceptors.tsConverter.type = timestamp agent.sources.file.interceptors.tsConverter.header = timestamp_str # 这里的pattern要和JSON中timestamp的格式完全匹配,比如带T的话用yyyy-MM-dd'T'HH:mm:ss agent.sources.file.interceptors.tsConverter.pattern = yyyy-MM-dd HH:mm:ss # SINKS CONFIGURATION(保持原配置不变) agent.sinks.hdfsSink.type = hdfs agent.sinks.hdfsSink.hdfs.path = /HADOOP/PATH/%Y/%m/%d/%H/ agent.sinks.hdfsSink.hdfs.filePrefix = common agent.sinks.hdfsSink.hdfs.fileSuffix = .json agent.sinks.hdfsSink.hdfs.rollInterval = 300 agent.sinks.hdfsSink.hdfs.rollSize = 5242880 agent.sinks.hdfsSink.hdfs.rollCount = 0 agent.sinks.hdfsSink.hdfs.maxOpenFiles = 2 agent.sinks.hdfsSink.hdfs.fileType = DataStream agent.sinks.hdfsSink.hdfs.callTimeout = 100000 agent.sinks.hdfsSink.hdfs.batchSize = 1000 agent.sinks.hdfsSink.channel = channel # CHANNELS CONFIGURATION(保持原配置不变) agent.channels.channel.type = memory agent.channels.channel.capacity = 10000 agent.channels.channel.transactionCapacity = 1000
关键注意事项
- 正则表达式匹配:要根据你JSON的实际格式调整正则,比如如果timestamp用单引号包裹,就把
"换成';如果字段名和值之间有多个空格,把\\s*保留(匹配任意空格)。 - 时间格式匹配:如果是字符串转时间戳,
pattern必须和JSON中timestamp的格式完全一致,否则转换会失败。 - 拦截器顺序:如果用多个拦截器,要保证先提取再转换,顺序不能搞反。
修改完配置后,重新启动Flume Agent,HDFS Sink就能正常读取事件Header中的timestamp来生成时间分区路径了。
内容的提问来源于stack exchange,提问作者frm
相关产品推荐
相关产品推荐

