【问题标题】:Apache flume custom interceptor - HDFS file in binary and strangeApache Flume 自定义拦截器 - 二进制和奇怪的 HDFS 文件
【发布时间】:2015-10-29 18:23:10
【问题描述】:

我对水槽拦截器概念相对较新,并且面临一个问题,即在应用拦截器之前,下沉的文件是普通文本文件,而在应用拦截器之后,一切都变得非常糟糕。

我的拦截器代码如下 -

package com.flume;

import org.apache.flume.*;
import org.apache.flume.interceptor.*;

import java.util.List;
import java.util.Map;
import java.util.ArrayList;
import java.io.UnsupportedEncodingException;
import java.net.InetAddress;
import java.net.UnknownHostException;

public class CustomHostInterceptor implements Interceptor {

    private String hostValue;
    private String hostHeader;

    public CustomHostInterceptor(String hostHeader){
        this.hostHeader = hostHeader;
    }

    @Override
    public void initialize() {
        // At interceptor start up
        try {
            hostValue =
                    InetAddress.getLocalHost().getHostName();
        } catch (UnknownHostException e) {
            throw new FlumeException("Cannot get Hostname", e);
        }
    }

    @Override
    public Event intercept(Event event) {

        // This is the event's body
        String body = new String(event.getBody());
        if(body.toLowerCase().contains("text")){
            try {
                event.setBody("hadoop".getBytes("UTF-8"));
            } catch (UnsupportedEncodingException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
        }
        // These are the event's headers
        Map<String, String> headers = event.getHeaders();

        // Enrich header with hostname
        headers.put(hostHeader, hostValue);

        // Let the enriched event go
        return event;
    }

    @Override
    public List<Event> intercept(List<Event> events) {

        List<Event> interceptedEvents =
                new ArrayList<Event>(events.size());
        for (Event event : events) {
            // Intercept any event
            Event interceptedEvent = intercept(event);
            interceptedEvents.add(interceptedEvent);
        }

        return interceptedEvents;
    }

    @Override
    public void close() {
        // At interceptor shutdown
    }

    public static class Builder
            implements Interceptor.Builder {

        private String hostHeader;

        @Override
        public void configure(Context context) {
            // Retrieve property from flume conf
            hostHeader = context.getString("hostHeader");
        }

        @Override
        public Interceptor build() {
            return new CustomHostInterceptor(hostHeader);
        }
    }
}

Flume 配置是 -

agent.sources=exec-source
agent.sinks=hdfs-sink
agent.channels=ch1

agent.sources.exec-source.type=exec
agent.sources.exec-source.command=tail -F /home/cloudera/Desktop/app.log
agent.sources.exec-source.interceptors = i1
agent.sources.exec-source.interceptors.i1.type = com.flume.CustomHostInterceptor$Builder
agent.sources.exec-source.interceptors.i1.hostHeader = hostname

agent.sinks.hdfs-sink.type=hdfs
agent.sinks.hdfs-sink.hdfs.path= hdfs://localhost:8020/bosch/flume/applogs
agent.sinks.hdfs-sink.hdfs.filePrefix=logs
agent.sinks.hdfs-sink.hdfs.rollInterval=60
agent.sinks.hdfs-sink.hdfs.rollSize=0

agent.channels.ch1.type=memory
agent.channels.ch1.capacity=1000

agent.sources.exec-source.channels=ch1
agent.sinks.hdfs-sink.channel=ch1

关于在 HDFS 中创建的文件上做猫 -

SEQ!org.apache.hadoop.io.LongWritable"org.apache.hadoop.io.BytesWritable���*q�CJv�/ESmP�ź
                                                                                           some textP�żc
                                                                                                           some more textP���K
                                                                                                                             textP��ߌangels and deamonsP��%�
          text bla blaP��1�angels and deamonsP��1�
                                                     testP��1�hmmmP��1�anything

有什么建议吗?

谢谢

【问题讨论】:

    标签: hadoop hadoop-streaming flume flume-ng


    【解决方案1】:

    Interceptor 看起来没什么问题。

    在您的 Flume 代理配置中。

    您没有指定此属性 (hdfs.fileType),因此将其作为默认序列文件

    尝试将此行添加到您的 HDFS SINK,如果可行,请告诉我。

    agent.sinks.hdfs-sink.hdfs.fileType=DataStream 
    

    【讨论】:

    • 确定一件事,我将把它标记为正确答案,它奏效了。然而,由于这是我第一次使用拦截器,我想了解它到底在做什么。我认为它会实时处理我的数据,并会实际检查包含“文本”的正文并将其替换为“hadoop”。这没有发生,有什么建议吗?
    • 您可以使用拦截器随意转换和丰富数据。使用此方法来执行 public Event intercept(Event event) { } Give Printout statements 并在此处根据需要调试和转换消息。
    • 如果您有更多问题,您可以发布单独的特定问题。此问题已正确回答,您可以将其标记为已回答,因为您说它有效。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-12
    相关资源
    最近更新 更多